Google L5 System Design โ€” Worked Examples

Ten prompts reported on 1point3acres for Google SWE loops, mostly 2026. Six worked in full, four condensed. Format note that matters more than any diagram: several reports describe the round as fully verbal โ€” no whiteboard, no code, with a near-silent interviewer. You set the pace, and the first ~10 minutes of requirements clarification is graded.

The L5 rubric & the opener script

What the round is actually scoring, and the 3-minute script that front-loads the graded part.

1. What the L5 bar looks like, signal by signal

SignalWhat clears the bar
RequirementsPins down QPS, latency SLO, consistency, region geometry, internal vs external before any box.
ArchitectureEvery box defended: what breaks without it, what it costs.
StorageNames the store and the access pattern that picked it (point lookup vs range scan vs append-heavy).
FailurePer-dependency: what happens when it dies, fail-open vs fail-closed, and which is right here.
ScaleHot-key and hot-shard analysis with a concrete mitigation.

The single most-cited rejection reason in the reports is jumping to architecture without clarifying. Not depth. Not novelty. That one.

2. The opener โ€” say this out loud in every mock

Five questions, ~3 minutes

1. Who calls this โ€” internal services or end users?
   (decides auth model, SLA, abuse surface)
2. Scale? Peak QPS, total entities, growth rate.
   (decides sharding and whether one box works)
3. Latency SLO at p99, and is it on the critical path?
   (decides sync vs async, cache vs no cache)
4. Consistency: can a reader see stale data, and for how long?
   (the single biggest architecture fork)
5. One region or global? Where are the writers?
   (decides replication topology and conflict handling)

Then: "I'll assume X, Y, Z unless you want to steer me." Stating assumptions out loud is what converts silence into a grade.

Then the shape of the hour

  • ~10 min clarify + capacity math out loud.
  • ~10 min high-level boxes only. Resist diving.
  • ~25 min deep dive โ€” the interviewer picks the box. Have the atomic sequence ready for the one that has a race in it.
  • ~10 min failure & trade-off. Volunteer this; don't wait to be asked.
  • ~5 min evolution: MVP โ†’ 10ร—, and where you'd re-architect.

Verbal-round adaptation: when there's no whiteboard, narrate structure explicitly โ€” "three tiers: edge, decision, storage; let me walk the write path first, then the read path." Without a diagram the interviewer can only grade the words, so the words have to carry the structure.

3. Reusable moves that score in almost every prompt

Say these unprompted

  • Idempotency key on every retryable write. Names the double-write race before it's asked.
  • Fail-open vs fail-closed, tied to the business consequence, not to taste.
  • Hot key mitigation: shard the key (key:{0..N}), or local pre-aggregate then flush.
  • Backpressure: what the caller sees when the queue is full. 429 with Retry-After beats silent drop.
  • Watermarks / bounded state for anything streaming, so state doesn't grow forever.

Capacity math you should have memorized

ThingNumber
Redis point op, same DC~0.5โ€“1 ms p99, ~100k ops/s/core
SSD random read~100 ยตs; ~50โ€“100k IOPS/device
Cross-region RTT (USโ†”EU)~80โ€“100 ms
Spanner cross-region commit~50โ€“150 ms
Kafka partition~10 MB/s sustained
1 day at 100k QPS~8.6 B events
Object storage~$20/TB/mo hot, ~$4/TB/mo cold

Distributed rate limiter full

L6 Infra/Platform onsite, SD round 1, 2026-01. "Design a large-scale distributed rate limiter service used internally by multiple teams." Interviewer pushed hard on trade-offs rather than novelty.

1. Requirements in one line

Functional

  • check(tenant, key, cost) โ†’ ALLOW | DENY called on the request path by any internal service.
  • Rules configurable per tenant, per API, per caller identity; hierarchical (org โ†’ team โ†’ key).
  • Multiple algorithms: fixed window, sliding window, token bucket with burst.
  • Observability: per-rule allow/deny counts, near-real-time.

Out of scope: authn (assume identity is already resolved), billing, WAF/bot detection.

Non-functional

  • 10โตโ€“10โถ decisions/sec aggregate, multi-tenant.
  • p99 < 5 ms added latency โ€” it's on every request's critical path.
  • Multi-region; a region partition must not take down callers.
  • Availability > accuracy: 99.99%+.

Core tension: a limiter that is perfectly accurate must be strongly consistent and therefore globally coordinated, which is exactly the thing you cannot afford on the critical path. Every decision below is buying accuracy back cheaply.

2. Core entities & schema

Rule            (config store, read-mostly, versioned)
  rule_id, tenant_id, scope        -- "api:/v1/search", "org:ads"
  algo                             -- TOKEN_BUCKET | SLIDING_WINDOW
  capacity, refill_per_sec, window_ms
  action_on_breach                 -- DENY | SHADOW | DEGRADE
  version, updated_at

Counter         (hot path, Redis / in-memory, TTL'd)
  key = {tenant}:{scope}:{subject}:{bucket_ts}
  value = tokens (float) | count (int)
  aux   = last_refill_ts
  TTL   = 2 ร— window            -- self-cleaning, no GC job

Decision log    (async, sampled, โ†’ Kafka โ†’ BigQuery-ish)
  ts, rule_id, subject, allowed, remaining, node_id

Why Redis for counters, not the SQL config store

The counter is a read-modify-write at request rate with a TTL and no need for durability โ€” losing a counter costs one over-admitted window, not money. That is exactly Redis's shape: single-threaded per shard so INCR/Lua are atomic without locks, and TTL eviction means no cleanup job. Putting it in SQL costs a disk fsync per request and a row-lock hot spot on popular tenants. Conversely rules are low-volume, must survive a full cache flush, and need audit history โ€” so they live in a replicated SQL store and are pushed to limiter nodes, never read on the request path.

3. API interfaces

Hot path (called by every service)

POST /v1/check
{ tenant, scope, subject, cost=1 }
โ†’ 200 { allowed: true,  remaining: 41, reset_ms: 730 }
โ†’ 200 { allowed: false, remaining: 0,  retry_after_ms: 730 }

// caller translates false โ†’ HTTP 429 + Retry-After
// batched variant amortizes RPC overhead:
POST /v1/check:batch  { checks: [...] }

Return remaining and reset_ms always โ€” callers surface them as X-RateLimit-* headers, which is what stops clients from hot-looping into you.

Control path (rare)

PUT   /v1/rules/{rule_id}     { algo, capacity, ... }
GET   /v1/rules?tenant=ads
POST  /v1/rules/{id}:shadow   // evaluate, never deny
DELETE /v1/rules/{rule_id}

GET /v1/stats?rule_id=..&window=5m
โ†’ { allowed, denied, p99_check_ms, top_subjects[] }

Shadow mode is the feature that gets you the offer. Nobody ships a new limit straight to DENY โ€” you run it in shadow, look at what would have been denied, then flip.

4. Architecture

Caller service (any team) Sidecar / SDK local token cache + batched flush Limiter service stateless, autoscaled runs Lua decision Redis Cluster sharded counters, TTL Rule store (SQL) versioned, pushed Kafka (decisions) sampled, async Stats / alerting top-N abusers Degraded path: Redis unreachable decide from local cache, then fail-open + alarm miss EVAL push 429 + Retry-After

Two paths: the common one never leaves the sidecar (local tokens pre-fetched in bulk); the miss path does one Redis Lua round-trip. Everything else โ€” rules, stats โ€” is off the critical path by construction.

5. Deep dive โ€” the atomic decision

Token bucket in one Lua script

-- KEYS[1] = bucket key
-- ARGV = now_ms, capacity, refill_per_ms, cost
local b = redis.call('HMGET', KEYS[1], 'tk', 'ts')
local tokens = tonumber(b[1]) or ARGV[2]
local ts     = tonumber(b[2]) or ARGV[1]

local delta  = math.max(0, ARGV[1] - ts)
tokens = math.min(ARGV[2], tokens + delta * ARGV[3])

local ok = tokens >= tonumber(ARGV[4])
if ok then tokens = tokens - ARGV[4] end

redis.call('HMSET', KEYS[1], 'tk', tokens, 'ts', ARGV[1])
redis.call('PEXPIRE', KEYS[1], 2 * ARGV[5])
return { ok and 1 or 0, tokens }

CAPACITY = burst allowance; refill_per_ms = steady-state rate. Two knobs, and being able to say which one a customer complaint maps to is the whole point of picking this algorithm.

Why it's correct, and where it isn't

Redis executes a Lua script atomically on the shard, so read-refill-decrement-write cannot interleave. That kills the classic race where two concurrent GET/SET pairs both see 1 token left and both allow.

Defend it: the script must be deterministic โ€” never call TIME inside it, always pass now_ms from the caller, or replicas diverge from the primary. That in turn means clock skew across limiter nodes leaks quota: a node 500 ms fast refills early. Bound it with NTP + reject now_ms more than ~1 s off the shard's own clock.

Why not sliding-window-log (a sorted set of timestamps per key)? Exact, but memory is O(requests in window) โ€” a 10k-QPS tenant on a 60 s window is 600k members in one key. Token bucket is O(1) memory and the accuracy loss is a bounded burst.

6. Deep dive โ€” hot keys and hierarchical limits

The hot-key problem

Rate limiter traffic is maximally skewed by design โ€” the abusive tenant is the one generating the most checks, and all of its checks hash to one Redis shard. That shard saturates and you have taken down limiting for every tenant on it.

Fix 1 โ€” key sharding:
  key = {tenant}:{scope}:{h(subject) % N}
  each sub-bucket gets capacity/N
  cost: burst accuracy degrades ~Nร—
        (a client can get N bursts)

Fix 2 โ€” local pre-allocation (preferred):
  sidecar leases 100 tokens at a time
  spends locally, refreshes at 20% left
  cost: up to (leases ร— lease_size) over-admit
        on abrupt traffic drop; leases TTL out

Hierarchical limits

Real rules nest: org 100k/s, team 10k/s, key 1k/s. Naรฏvely you evaluate all three, which triples Redis load and creates a partial-decrement bug โ€” decrement org, then deny at team, and org's tokens are gone.

Fix: evaluate all levels in one Lua script, cheapest/narrowest first, and only commit decrements if every level passes. Same-slot placement via hash tags: {tenant}:org, {tenant}:team:x โ€” the braces force Redis Cluster to co-locate them so a multi-key script is legal.

Belt-and-suspenders: if levels genuinely cannot be co-located, decrement narrowest-first and refund on a later-level deny. Refunds are idempotent (add back cost, cap at capacity) so a lost refund costs one wasted token, not a stuck bucket.

7. Deep dive โ€” multi-region

Don't replicate counters synchronously

A globally-accurate limit needs consensus per decision: ~100 ms cross-region RTT against a 5 ms budget. Non-starter. Three options, in order of what you should actually propose:

ModelAccuracyCost
Per-region quota split (limit/R per region)Never over-admits globally; under-admits on skewFree
Split + periodic rebalance (every ~10 s, redistribute unused)GoodOne background job
Global consensus (Spanner counter)ExactBlows the latency SLO

What to say

"I'd start with static per-region split because it's strictly safe โ€” the global limit is never exceeded. The failure mode is that a region with 90% of traffic gets throttled while others sit idle, so I'd add an async rebalancer that redistributes unused quota on a 10-second loop. That converts a correctness problem into a utilization problem, which I'd rather have."

Defend it: when a region is partitioned, its quota share is stranded. The rebalancer must treat a silent region as still holding its share (don't redistribute), or you over-admit globally the moment it comes back. Conservative on partition, aggressive on healthy โ€” that asymmetry is the whole trick.

8. Follow-ups โ€” answers to have ready

Redis dies. Now what?

Fail-open, loudly. This is an internal limiter protecting against accidental overload, not a security control โ€” failing closed converts a cache outage into a total outage of every calling service, which is strictly worse than briefly unlimited traffic. Concretely: sidecar serves from its last local lease, then admits everything, emits a limiter_degraded metric, and pages. The exception is any rule tagged security (login attempts, password reset) โ€” those fail closed, and the rule schema carries that flag so the decision is made at config time, not during the incident.

Two limiter nodes both serve the same tenant. Do they double-admit?

No, because neither node holds state โ€” the counter lives on one Redis shard and the Lua script is atomic there. The only double-admit is via the local-lease optimization, and it's bounded: at most nodes ร— lease_size extra requests, which you tune by shrinking leases for tight limits and growing them for loose ones.

A tenant complains they're getting 429s below their stated limit.

Three usual causes, in order of likelihood. (1) Fixed-window boundary effect โ€” they sent 2ร— limit across a window edge on a previous design; sliding/token-bucket fixes it. (2) Key sharding split their capacity N ways and their traffic isn't uniformly distributed across subjects. (3) Clock skew on one limiter node making refill run slow. The remaining/reset_ms in every response plus the sampled decision log is what lets you answer this in minutes rather than guessing.

How do you roll out a new limit without breaking a team?

Shadow mode: evaluate, log, never deny. Look at the would-be-denied distribution for a week, notify the top offenders, then flip to enforce at a generous capacity and ratchet down. Rules are versioned so rollback is a config push, not a deploy.

Why not just do this at the load balancer?

You should, for the crude global limits โ€” L7 LB rate limiting is cheaper and closer to the client. But it can't see application-level identity (which tenant, which API key, which cost weight), can't do hierarchical limits, and can't be reconfigured per-team without touching shared infra. Layered: LB for the blunt DDoS ceiling, this service for the semantic limits.

Requests have different costs. Does that break the model?

No โ€” cost is already a parameter, so an expensive search costs 10 tokens and a health check costs 0. The subtlety is that you don't know the true cost until after execution. Reserve an estimated cost up front, then reconcile with an async adjust call. Over-reserving is safe; under-reserving lets one expensive request slip, which is acceptable.

9. Numbers to drop

Load

  • 10โถ checks/sec peak aggregate.
  • With sidecar leases at 90% local hit rate โ†’ 10โต Redis ops/sec reaching the cluster.
  • Redis shard sustains ~100k ops/s โ†’ ~4 shards + replicas covers it with headroom. That's a startlingly small cluster, and saying so is the point.
  • Limiter service: ~5k checks/s/pod โ†’ ~200 pods, stateless, trivially autoscaled.

Storage & latency

  • Counter entry ~100 B; 50 M active keys โ†’ ~5 GB. Fits in memory on one shard set; TTL keeps it flat.
  • Rules: 10โต rules ร— 1 KB = 100 MB. Fits in every node's memory โ€” that's why you push, not fetch.
  • Budget: local hit <0.1 ms; Redis round-trip ~1 ms p99; total p99 under 5 ms with room to spare.
  • Decision log at 1% sampling โ†’ 10k events/s โ†’ ~1 Kafka partition-pair. Full logging would be 100ร— and isn't worth it.

10. 30-second recap

Callers hit a sidecar that holds a small lease of tokens locally, so ~90% of checks never leave the process. On a miss it calls a stateless limiter service, which runs one atomic Lua script against a sharded Redis Cluster โ€” read tokens, refill by elapsed time, decrement if sufficient, all in one shot, so there's no read-modify-write race. Rules live in a versioned SQL store and are pushed to nodes, never fetched on the hot path. Multi-region is a static quota split with an async rebalancer, because global consensus per decision blows the 5 ms budget. If Redis dies we fail open and page, since this protects against accidental overload rather than attackers โ€” except for rules explicitly tagged security, which fail closed. The bounded inaccuracy is nodes ร— lease_size extra admits, and I'd shrink leases for tight limits to trade Redis load for precision.

Global real-time notification system full

L6 onsite SD round 2, 2026-01. "Design a global real-time notification system โ€” push + email + SMS." Stated constraints: 10โธ+ users, multi-channel fallback, latency SLA, dedup and idempotency. The poster's read: this round tested system-owner thinking more than the infra round did.

1. Requirements in one line

Functional

  • Producer services submit a notification event; the platform decides channel, renders, delivers, retries.
  • Channels: mobile push (APNs/FCM), email, SMS, in-app inbox. Fallback chain per event type.
  • User preferences and quiet hours honored; unsubscribe is legally binding for email/SMS.
  • Fan-out: one event may target one user, a segment, or all users (announcement).
  • Exactly-once user-visible delivery: at-least-once transport + dedup at the edge.

Out of scope: content authoring UI, ML send-time optimization, marketing campaign scheduling.

Non-functional

  • 10โธ users; steady ~50k notifications/s, burst 1 M/s on a broadcast.
  • Transactional (2FA, security alert): p99 < 5 s end-to-end. Marketing: minutes are fine.
  • Durability > latency for transactional โ€” never silently drop.
  • Multi-region, survives loss of one region and of any single third-party vendor.

Core tension: the burst is 20ร— steady and the third-party channels (APNs, an SMS vendor) are both rate-limited and unreliable. So the system's real job is absorbing bursts and metabolizing partner failure, not sending messages.

2. Core entities & schema

NotificationEvent      -- what the producer asked for (immutable)
  event_id (idempotency key, producer-supplied)
  type            -- ORDER_SHIPPED | SECURITY_ALERT | PROMO_X
  audience        -- {user_id} | {segment_id} | ALL
  payload         -- template vars, NOT rendered text
  priority        -- TRANSACTIONAL | STANDARD | BULK
  created_at, ttl_sec

Delivery               -- one row per (event, user, channel). The unit of work.
  delivery_id = hash(event_id, user_id, channel)   <-- dedup key
  state       -- PENDING|SENT|DELIVERED|FAILED|SUPPRESSED
  attempt, next_attempt_at, last_error
  vendor_msg_id

UserPrefs              -- read-heavy, cached hard
  user_id, channel_enabled{}, quiet_hours, locale, tz
  device_tokens[]  (token, platform, last_seen)
  suppression: unsubscribed_types[], hard_bounce, complaint

Why split Event from Delivery โ€” and why this store

Modelling only "notification" is the mistake that sinks this round. One event fans out to N users ร— M channels, each with its own retry clock and terminal state; without a Delivery row you cannot answer "did user X get it?" or retry one channel without re-sending the others. It also makes dedup a primary-key property: delivery_id is a deterministic hash, so a duplicated event replays into the same row instead of a second send.

Delivery is high-volume, append-then-update-once, always accessed by key or by (state, next_attempt_at) โ€” that's a wide-column store (Bigtable/Cassandra) with a TTL, not a relational DB; there are no joins and no transactions across users. UserPrefs is small (10โธ ร— ~1 KB = 100 GB), read on every delivery, and must be correct for unsubscribes, so: replicated SQL as source of truth, aggressively cached, with suppression checked again at send time.

3. API interfaces

Producer (write path)

POST /v1/notify
Idempotency-Key: {event_id}
{ type, audience: {user_id | segment_id},
  payload: {...}, priority, ttl_sec }
โ†’ 202 { event_id, status: "ACCEPTED" }

POST /v1/notify:broadcast
{ type, segment_id, payload, rate_limit_per_sec }
โ†’ 202 { event_id, estimated_recipients }

202, never 200. Accepting into a durable log and returning immediately is what decouples the producer from APNs being down. A producer that blocks on delivery has coupled its own availability to a third party.

Consumer & ops (read path)

GET  /v1/inbox?user_id=&cursor=   // in-app feed
POST /v1/deliveries/{id}:ack     // client read receipt
GET  /v1/events/{event_id}       // fan-out status rollup
     โ†’ { total, sent, delivered, failed, suppressed }

POST /v1/webhooks/vendor/{vendor}  // bounce/delivery callbacks
PUT  /v1/prefs/{user_id}
POST /v1/unsubscribe   { token }   // one-click, no login

4. Architecture

Producers order, auth, promo Ingest API idempotency check + 202 Event log Kafka, partitioned by priority Fan-out workers expand segment, apply prefs, write Delivery rows UserPrefs cache + suppression list Delivery store wide-column, TTL Per-channel queues (push / email / sms / inbox) separate rate budgets, separate blast radius Push sender APNs / FCM Email sender vendor A / B SMS sender vendor A / B In-app inbox always written Retry scheduler + DLQ exp backoff, TTL expiry, channel fallback Vendor webhooks โ†’ state machine delivered / bounce / complaint โ†’ suppression

Green is the durable happy path: accept โ†’ log โ†’ fan-out โ†’ per-channel queue. Amber is the recovery loop that everything eventually falls into. The in-app inbox is written unconditionally, so there is always one channel that cannot fail.

5. Deep dive โ€” fan-out and the broadcast burst

Two fan-out modes

Targeted (99% of events, 1..1000 users):
  fan-out inline in the worker
  write Delivery rows in a batch
  enqueue per channel

Broadcast (segment or ALL, 10^8 users):
  DO NOT expand in one worker
  1. materialize segment -> sharded user-id ranges
  2. emit N "fan-out chunk" tasks (chunk = 10k users)
  3. each chunk expands independently, resumable
     via (event_id, chunk_id) checkpoint
  4. token-bucket the emit rate per channel

A broadcast at 10โธ users writing 10โธ Delivery rows is ~100 GB and hours of vendor throughput. It must be a paced background job, not a spike.

Why chunking, and where it still hurts

Chunk-level checkpoints make fan-out resumable and idempotent: a worker crash replays one 10k chunk, and the deterministic delivery_id means the replayed rows collide with the originals instead of duplicating. Without this, a crash 80% through a 10โธ fan-out leaves you with no safe action.

Defend it โ€” the priority inversion. A broadcast dumped into the same queue as 2FA codes will bury them behind hours of backlog, and a 2FA code delivered in 40 minutes is worse than useless. Fix: separate physical queues per priority class, not just a priority field. Transactional gets its own partitions, its own worker pool, and a reserved share of vendor rate budget that bulk traffic can never borrow. This is the single most likely follow-up in this round; volunteer it.

6. Deep dive โ€” exactly-once as the user perceives it

The chain of idempotency

1. Producer supplies event_id (Idempotency-Key).
   Ingest does INSERT ... IF NOT EXISTS.
   Retry of the same HTTP call -> same event, no dup.

2. delivery_id = hash(event_id, user_id, channel)
   Fan-out replay writes the same row. Dedup by PK.

3. Send is guarded by a conditional state transition:
     UPDATE Delivery
        SET state='SENDING', attempt=attempt+1
      WHERE delivery_id=? AND state IN ('PENDING','FAILED')
   Only the winner calls the vendor.

4. Vendor call carries its own idempotency token
   where supported (SES, Twilio); where not, accept
   at-least-once and dedup on the DEVICE by event_id.

Why true exactly-once is impossible, and the honest answer

Step 3 has an unavoidable gap: you set SENDING, call the vendor, and crash before recording the result. You cannot know whether the vendor sent it. Any system that claims exactly-once across a third-party boundary is lying.

So decide per channel, by consequence. Push and in-app: retry (duplicate push is mildly annoying, and the client dedups on event_id anyway). Email: retry (dup email is tolerable). SMS: do not blind-retry โ€” it costs real money and duplicate 2FA codes confuse users and can invalidate the first code. For SMS, reconcile against the vendor's delivery-status API before a retry, and cap at one reconciliation attempt.

Belt-and-suspenders: the client-side dedup on event_id with a 24 h local cache is what actually delivers the user-visible guarantee. Server-side you only ever promise at-least-once.

7. Deep dive โ€” retries, fallback, and vendor failure

Retry policy that doesn't amplify an outage

backoff = min(2^attempt * 1s, 15m) * jitter(0.5..1.5)
max_attempts = 5, then DLQ
hard stop at event.ttl_sec  -- a 2FA code has ttl 300s;
                            -- retrying at t+10m is harmful

Circuit breaker per (channel, vendor):
  error_rate > 50% over 30s  -> OPEN
  OPEN: stop calling, park in queue, alarm
  half-open probe every 10s

Jitter is not optional. Without it, a vendor recovering from a 60 s outage gets every parked message at the same instant and immediately falls over again.

Channel fallback โ€” and why it's usually wrong

The prompt asks for multi-channel fallback, so propose it, then bound it. Fallback chain lives on the event type: SECURITY_ALERT: [push, sms, email], PROMO: [push] (never escalate marketing to SMS โ€” it costs money and generates complaints).

Defend it: escalating on "not delivered within N seconds" is dangerous, because push delivery receipts are unreliable and slow โ€” you'll double-notify constantly. Only escalate on a hard signal: no registered device token, token rejected by APNs, or circuit breaker open for that channel. Timeout-based escalation should have a generous floor (minutes, not seconds) and must be off for anything bulk.

Vendor redundancy: two vendors per paid channel, weighted routing, automatic weight shift when a breaker opens. This is the concrete answer to "SMS vendor dies" and it's cheap to say.

8. Follow-ups โ€” answers to have ready

APNs goes down for an hour. What does the user see?

Nothing missing, because the in-app inbox row is written at fan-out time regardless of push outcome โ€” open the app and it's there. Push deliveries park behind an open circuit breaker; the retry scheduler holds them until TTL. When APNs recovers, the half-open probe reopens the valve and the backlog drains with jitter. Anything whose TTL expired in the meantime is dropped to DLQ and counted โ€” and for security alerts specifically, the fallback chain escalates to SMS on breaker-open rather than waiting.

Kafka loses a partition mid-fan-out.

Events are durably committed before the 202, and consumer offsets only advance after Delivery rows are written, so a partition failover replays from the last committed offset. Replay is safe precisely because delivery_id is deterministic โ€” re-processed chunks collide on the primary key instead of duplicating. The visible cost is latency during failover, not lost notifications. Retention of 7 days also means we can rewind and replay a window if we ship a fan-out bug.

A hot user โ€” say a celebrity account triggering millions of follower notifications.

That's the broadcast path, not the targeted path, and the trigger should classify it as such above a threshold (say >10k recipients). Practically: materialize the follower list into chunks, pace the emit with a token bucket, and โ€” importantly โ€” collapse. If the same producer generates 50 events for one user in a minute, digest them into one notification rather than sending 50. Collapse keys are the standard tool (collapse_key = (user_id, type), keep latest), and APNs/FCM support it natively.

User unsubscribes. How fast does that take effect, and what if a fan-out is in flight?

Suppression is checked twice: once at fan-out (cheap, from cache) and again in the sender immediately before the vendor call (authoritative read). The second check exists exactly for the in-flight case โ€” a broadcast queued an hour ago must not send to someone who unsubscribed ten minutes ago. For email/SMS this isn't a nicety; it's CAN-SPAM/TCPA exposure, and saying that out loud signals you've shipped this before.

How do you handle quiet hours across time zones?

Store the user's tz and quiet window in prefs; at fan-out, transactional messages ignore quiet hours entirely, and non-transactional ones get next_attempt_at set to the end of the quiet window rather than being dropped. Consequence: a global broadcast doesn't go out at one instant, it rolls around the planet over ~24 h, which you should state explicitly because it changes the "burst" math from 1 M/s to something far gentler.

How do you know delivery is actually working?

Per (type, channel, vendor, region): submitted โ†’ sent โ†’ delivered โ†’ opened funnel, with alerting on ratio shifts, not absolute counts. The signal that catches real incidents is delivered/sent dropping while sent stays flat โ€” that means the vendor is accepting and silently dropping, which no error rate would show. Plus a synthetic canary: send to a handful of owned test accounts on every channel every minute and alert on the round-trip.

Region loss.

Producers are geo-routed and events are replicated cross-region in the log. Fan-out and senders are active-active with the Delivery store as the arbiter โ€” the conditional PENDING โ†’ SENDING transition prevents two regions from sending the same delivery, assuming the store does per-key linearizable writes (Spanner or a quorum config). If it can't, accept the duplicate and rely on client-side dedup; say which one you're assuming.

9. Numbers to drop

Throughput

  • Steady 50k notifications/s โ†’ ~4 B/day. At ~1 KB of Delivery row each, ~4 TB/day before TTL.
  • 30-day TTL on Delivery โ†’ ~120 TB. Wide-column with compression, ~20 nodes. Cheap.
  • Kafka: 50k msg/s ร— 1 KB = 50 MB/s โ†’ ~10 partitions at 10 MB/s, call it 32 for headroom and priority isolation.
  • Broadcast to 10โธ: at a paced 100k sends/s that's ~17 minutes per channel โ€” quote this number, it reframes "real-time" honestly.

Latency budget (transactional, 5 s p99)

  • Ingest + commit to log: ~20 ms
  • Log โ†’ fan-out worker pickup: ~100 ms
  • Prefs lookup (cached): ~2 ms
  • Delivery write + queue: ~30 ms
  • Vendor call (APNs): ~200 ms p99
  • Device receipt: ~1 s, outside our control
  • Slack: ~3.6 s โ€” which is what pays for exactly one retry inside the SLA.

10. 30-second recap

Producers POST an event with an idempotency key and get a 202 as soon as it's durable in a partitioned log โ€” nothing blocks on a third party. Fan-out workers expand the audience, apply prefs and suppression, and write one Delivery row per (event, user, channel) with a deterministic ID, which is what makes the whole pipeline replay-safe. Deliveries go into per-channel, per-priority physical queues so a 10โธ-user broadcast can't bury a 2FA code. Senders take a conditional PENDINGโ†’SENDING transition before calling the vendor, so only one worker sends; retries use exponential backoff with jitter behind a per-vendor circuit breaker, and hard-stop at the event TTL. Fallback escalates on hard signals only โ€” no device token, or breaker open โ€” never on a soft delivery timeout. The honest limit is that exactly-once across a vendor boundary is impossible, so we promise at-least-once server-side and dedup on the device by event ID, and we don't blind-retry SMS because it costs money.

Anomaly / metrics monitoring system full

L6 GCP phone-screen SD, 2026-05 (reported as "SD: ๅผ‚ๅธธ็›‘ๆต‹็ณป็ปŸ"). The question bank frames it as deliberately under-specified โ€” clarify the domain first (metric stream? user behavior? infra alerts?). The expected reading: a metrics platform that ingests, stores, detects, alerts, and visualizes time series.

1. Requirements in one line

Functional

  • Ingest multi-dimensional time series: metric{service,region,host,env,segment} โ†’ value@ts.
  • Query: aggregate over label selectors and time ranges, for dashboards and detectors.
  • Detection, at least two mechanisms: static threshold and dynamic baseline (seasonality / change-point).
  • Alerting: routing by team, dedup, grouping, silencing, escalation, multiple channels.
  • Visualization: dashboards, plus the ability to slice by dimension after an alert fires.

Out of scope: log search, distributed tracing, incident management workflow (page a PagerDuty-equivalent, don't build it).

Non-functional

  • 10โท active series, 10โท datapoints/s ingest.
  • Detection latency < 60 s from event to page for critical metrics.
  • Query p99 < 1 s for a 24 h dashboard panel.
  • The monitoring system must not depend on the systems it monitors, and must degrade before it lies.

Core tension: alert quality is a precision/recall trade-off, and it is a product decision, not a modelling one. A detector that never misses will page so often that humans stop reading, at which point recall is effectively zero. Everything below is about being deliberately less sensitive in exchange for being believed.

2. Core entities & schema

Series          -- identity is the full label set
  series_id = hash(metric_name, sorted(labels))
  metric_name, labels{}, type (gauge|counter|histogram)
  first_seen, last_seen

Sample          -- the hot, huge one
  series_id, ts, value
  stored COLUMNAR + delta-of-delta on ts,
  XOR/Gorilla on value  -> ~1.4 bytes/sample

Detector
  detector_id, selector (PromQL-ish), algo,
  params{threshold | sensitivity | season_len},
  for_duration,        -- must hold N min before firing
  severity, route_to (team), runbook_url

AlertInstance   -- the dedup unit
  fingerprint = hash(detector_id, grouping_labels)
  state -- PENDING|FIRING|RESOLVED|SILENCED
  started_at, last_eval, value, notification_count

Why a purpose-built TSDB and not "just Bigtable"

Three properties make time series a special case, and naming them is the point of this section. (1) Writes are append-only and time-ordered, so you never update a sample โ€” that kills the need for compaction-heavy general stores. (2) Adjacent values are highly correlated, so Gorilla-style XOR + delta-of-delta compression gets you ~1.4 bytes per sample versus ~16 raw; at 10โท points/s that's the difference between 500 TB/year and 4 PB/year. (3) Queries are range scans over one series, so the physical layout must be series-major, time-minor โ€” the exact opposite of what you'd get by keying on timestamp first (which also creates a write hot spot on "now").

Concretely: shard by hash(series_id) so writes spread evenly and one series' history is contiguous. An inverted index maps label=value โ†’ [series_id] so a selector resolves to a series set before touching sample data.

3. API interfaces

Write path

POST /v1/ingest            (protobuf, batched, gzip)
{ samples: [{ labels{}, ts, value }, ... ] }
โ†’ 204, or 429 with Retry-After when over quota

// agents push every 10s; server does NOT pull.
// per-tenant series-cardinality quota enforced HERE.

Push, not pull, at this scale โ€” pull requires the monitoring system to hold a service-discovery view of 10โท targets and creates a fan-in storm. Pull's advantage (you know when a target is missing) is recovered by a staleness detector.

Read / control path

GET /v1/query_range
  ?selector=http_errors{service="ads",env="prod"}
  &start=&end=&step=60s&agg=sum by (region)
โ†’ { series: [{labels, points:[[ts,v]...]}] }

PUT  /v1/detectors/{id}
POST /v1/silences  { matcher, until, reason }
GET  /v1/alerts?state=FIRING
POST /v1/detectors/{id}:backtest  { window: "30d" }

Backtest is the feature that makes detectors shippable โ€” replay a candidate detector over 30 days of history and show how many times it would have paged. Nobody should tune a threshold in production.

4. Architecture

Agents 10s scrape Ingest tier validate, quota, cardinality guard WAL / Kafka partition by series_id Streaming detect sliding window, <60s path TSDB (hot) in-mem + SSD, 2h blocks Object store (cold) downsampled 5m/1h Alert manager dedup, group, silence, escalate Query engine hot + cold merge Batch detectors seasonal baselines Notifiers page / IM / email Dashboards slice by label Meta-monitoring (separate stack) watches this system; dead-man's switch compact

Two detection paths deliberately: a streaming one on the raw feed for the <60 s threshold alerts, and a batch one reading the TSDB for seasonal baselines that need hours of context. Both converge on one alert manager, which is the only thing allowed to page.

5. Deep dive โ€” the two detectors

Static threshold (streaming)

for each eval tick (15s):
  v = agg(selector, last 5m)
  if compare(v, threshold):
     pending[fp] = pending[fp] or now
     if now - pending[fp] >= for_duration:   # e.g. 5m
        emit FIRING
  else:
     clear pending[fp]; emit RESOLVED if firing

for_duration is the anti-flap knob, and it's the whole reason threshold alerting is usable. A metric crossing a line for 15 s is noise; crossing it for 5 minutes is an incident. Cost: you've added for_duration to your detection latency, so critical detectors run a short duration on a tight threshold and a long duration on a loose one.

Dynamic baseline (batch)

Seasonal decomposition per series:
  baseline(t) = median over last K weeks
                at same (weekday, time-of-day)
  band(t)     = baseline(t) ยฑ k * MAD(t)
  fire when actual outside band for for_duration

Change point (fast, cheap):
  compare rolling mean/var of last 5m vs prior 60m
  robust z = (m1 - m0) / (1.4826 * MAD_0)
  |z| > 4  -> candidate

Median + MAD, not mean + stddev. Mean and standard deviation are wrecked by the very outliers you're trying to detect โ€” one past incident inflates the band and hides the next one. MAD is robust to ~50% contamination, which is the difference between a detector that degrades gracefully and one that silently stops working after its first real incident.

Defend it โ€” why not "just use ML"

Reach for a learned model only when you can state what it buys. Here it buys multivariate correlation (error rate up and latency up and only in one region) at the cost of: training pipelines, a cold-start period per new service, opacity at 3 a.m. when the on-call asks why it fired, and a whole new failure mode where the model drifts and nobody notices. Seasonal median + MAD covers most of the value at a fraction of the operational cost. Propose ML as a ranking layer over already-fired alerts (which of these 40 alerts is the root cause), not as the trigger โ€” that framing is far more defensible than "add an LSTM."

6. Deep dive โ€” cardinality, the actual failure mode

How these systems really die

Not from datapoint volume โ€” from series count. Every distinct label combination is a new series with its own index entry and memory-resident head block. One engineer adds user_id or request_id or a raw URL as a label and cardinality goes from 10โท to 10โน overnight, the index blows out of RAM, and the whole cluster falls over. This takes down monitoring for everyone, during whatever incident prompted them to add the label.

Defenses, in order:
1. Per-tenant series quota, enforced at INGEST.
   Over quota -> reject new series, keep existing.
   (Never reject existing: partial data is worse
    than no new dimensions.)
2. Label-value cardinality cap per label name.
   >1000 distinct values -> auto-drop label, warn.
3. Cardinality explosion detector on the metadata
   itself: d(series_count)/dt alert.
4. Chargeback: teams see their series cost.

The other structural risks

Hot shard: hashing by series_id spreads writes evenly by construction, but a query for sum by (region) over a million series fans out to every shard. Mitigate with pre-computed recording rules โ€” materialize the common aggregations at write time so dashboards read one cheap series instead of a million expensive ones.

Alert storm: one bad deploy fires 400 detectors at once and pages a human 400 times. The alert manager must group by a shared label set (e.g. all alerts for service=ads,region=us-east become one notification with 400 lines) and apply an inhibition rule โ€” if ServiceDown is firing, suppress every dependent HighLatency alert for that service. Without inhibition rules the system is technically correct and operationally useless.

Silent failure: the worst outcome isn't a false page, it's no page because ingest stopped. Hence the dead-man's switch below.

7. Deep dive โ€” not lying, and not depending on yourself

Staleness > absence

An absent series and a healthy series look identical to a threshold detector: neither crosses the line. So every detector needs an explicit staleness clause โ€” absent_over_time(selector[5m]) fires its own alert. Otherwise a crashed agent reads as "all clear," which is the single most dangerous bug a monitoring system can have.

Dead-man's switch: a synthetic detector that fires continuously and routes to an external service which pages you when the heartbeat stops. It's the only construct that catches "the monitoring system itself is down," and it costs about ten lines.

Circular dependency

If the monitoring stack runs on the same Kubernetes control plane, the same service mesh, and the same object store as production, then a production outage blinds you exactly when you need sight. State this explicitly: the meta-monitoring stack runs in a separate failure domain โ€” different region, different account, minimal dependencies, and its notification path does not traverse the primary network. It monitors ~50 series about the main system rather than 10โท, so it can be small and boring.

Degrade before lying: under ingest overload, shed low-priority tenants and keep critical ones at full fidelity, rather than uniformly downsampling everyone. Uniform degradation quietly reduces detection sensitivity across the board with no signal that it happened.

8. Follow-ups โ€” answers to have ready

Kafka/WAL backs up. What gives?

Ingest applies backpressure with 429s and agents buffer locally (bounded, ~15 min, then drop oldest). Critically, the streaming detector reads from the same log, so a backup delays detection โ€” that's the real cost, not lost data. So: prioritized partitions, with detector-relevant series on a reserved partition set that bulk dashboard-only metrics can't saturate. And the lag itself is a first-class alert on the meta stack.

Someone asks for a 2-year retention dashboard.

Not at raw resolution โ€” 10โท points/s for 2 years is ~900 TB compressed and no one needs 10-second granularity from 18 months ago. Downsample on compaction: raw for 15 days, 5-minute rollups for 90 days, 1-hour rollups for 2 years, each with min/max/sum/count so you can still reconstruct averages and see spikes. The query engine picks resolution from the requested range automatically. Say the trade-off out loud: you permanently lose the ability to see a 30-second spike from last year.

How do you stop alert fatigue?

Treat it as the primary metric of the system, not a side effect. Track per-detector precision (fired โ†’ was there an actual incident?) via a required post-alert disposition, and auto-quarantine detectors below ~30% precision into a non-paging channel until the owner fixes them. Combine with grouping, inhibition, and for_duration. The organizational move matters as much as the technical one: pages should have a runbook link, and a page with no runbook shouldn't be allowed to page.

Two regions both evaluate the same detector. Double page?

Detector evaluation is sharded by fingerprint with a lease, so exactly one evaluator owns a detector at a time; on lease expiry another takes over. Even if two do evaluate during a partition, the alert manager dedups on fingerprint and the notification has its own dedup window, so the human sees one page. Duplicate evaluation is cheap and safe; duplicate notification is what you must prevent, and you prevent it at the last hop.

A metric is anomalous but it's Black Friday.

This is why the baseline is seasonal by (weekday, time-of-day) rather than a flat threshold โ€” but a once-a-year event isn't in the seasonal profile. Two mechanisms: scheduled silences tied to known events, and a relative comparison against a control group (this region vs all regions, this version vs previous). Relative detectors survive traffic-shape changes that absolute ones can't, which is a strong point to volunteer unprompted.

How would this evolve at 10ร—?

The ingest and TSDB tiers scale horizontally by series hash, so 10ร— is more shards โ€” boring, which is the point. The parts that don't scale linearly are the inverted index (a selector matching millions of series gets slow, so push more work into recording rules) and the human alert budget, which doesn't scale at all. At 10ร— the re-architecture is on the alerting side: correlate and rank rather than route more pages.

9. Numbers to drop

Ingest & storage

  • 10โท series ร— 1 sample / 10 s = 10โถ datapoints/s steady; size for 10โท peak.
  • At ~1.4 bytes/sample compressed: 10โถ/s โ†’ ~120 GB/day raw-resolution.
  • 15-day raw retention โ†’ ~1.8 TB. Genuinely small โ€” the compression is the headline number.
  • Index: ~10โท series ร— ~200 B = 2 GB of index in RAM, which is why cardinality (not volume) is the binding constraint.
  • Ingest node handles ~200k samples/s โ†’ ~50 nodes at peak.

Detection latency budget (60 s)

  • Agent scrape interval: 10 s (worst-case age at emit)
  • Agent โ†’ ingest โ†’ log commit: ~2 s
  • Streaming detector eval tick: 15 s
  • Alert manager group_wait: 10 s (batches the storm)
  • Notifier โ†’ page: ~3 s
  • ~40 s, leaving ~20 s of slack โ€” which is exactly the budget that for_duration eats, so critical detectors run for_duration=0 on high thresholds.

10. 30-second recap

Agents push batched samples to a stateless ingest tier that enforces per-tenant series-cardinality quotas โ€” cardinality, not volume, is how these systems actually die. Samples land in a partitioned log, which feeds two consumers: a streaming detector for sub-60-second threshold alerts, and a TSDB sharded by series ID storing Gorilla-compressed columnar blocks at about 1.4 bytes a sample. Batch detectors read the TSDB to compute seasonal baselines using median and MAD rather than mean and stddev, because the outliers you're hunting would otherwise inflate your own band. Everything funnels into one alert manager that dedups on fingerprint, groups by shared labels, applies inhibition rules so a service-down alert suppresses its dependents, and honors silences. Two things I'd insist on: an explicit staleness detector, because absent data and healthy data look identical to a threshold, and a dead-man's switch on a separate stack in a separate failure domain โ€” otherwise the monitoring system goes blind precisely when production breaks.

Find-My-Device full

L6 GCP onsite SD, 2026-05. "่ฎพ่ฎก find my iphone๏ผŒ้œ€่ฆๅˆ็†่€ƒ่™‘้š็งๆ€ง๏ผŒๅฎ‰ๅ…จๆ€งๅ’Œๆ•ˆ็އ." The prompt names privacy, security, and efficiency as the graded axes โ€” that's unusual and it means a correct-but-generic location-tracking design fails this round.

1. Requirements in one line

Functional

  • Owner sees near-real-time location of their devices from web or another device.
  • Lost mode: play sound, mark as lost with a message, remote lock, remote wipe.
  • Offline devices: report position when they come back online; optionally crowd-sourced reporting via nearby devices.
  • Family sharing: explicitly authorized members can view, with revocation.
  • Audit log the owner can read: who queried my location, when.

Out of scope: the map rendering stack, device activation lock enforcement, carrier integration.

Non-functional

  • 10โธ devices; peak 10โถ location updates/sec.
  • Query QPS is tiny by comparison โ€” people check rarely. Write-dominated by ~1000:1.
  • Battery: location reporting must cost single-digit mW average.
  • The server must not be able to read device locations. Treat that as a hard requirement, not a nice-to-have.

Core tension: a normal design puts plaintext locations in a database so the service can index and query them efficiently. But a database of where 10โธ people are is the highest-value breach target you could build, and an insider-abuse vector. The entire design is about giving that up โ€” and then paying for it in query complexity and lost server-side features.

2. Core entities & schema

Device
  device_id, owner_id, model, enrolled_at
  public_key (device keypair, private key never leaves device)
  status -- NORMAL | LOST | WIPED

LocationReport          -- what the server actually stores
  report_key = H(rotating_public_key)     <-- opaque, no device_id!
  ciphertext              -- E(location, timestamp) under device key
  received_at             -- server clock, for TTL only
  reporter_hint           -- coarse, for abuse throttling only
  TTL = 7 days
  // server CANNOT link report_key -> device_id -> user

Command                 -- the write-back channel
  command_id, device_id, type (SOUND|LOCK|WIPE|MARK_LOST)
  issued_by, issued_at, nonce
  signature (signed by owner's key)
  state -- PENDING | DELIVERED | ACKED | EXPIRED

AccessGrant
  grantor_id, grantee_id, device_id, expires_at, revoked_at
AuditEntry
  device_id, actor_id, action, ts, ip_coarse

Why the location store is a dumb blob store, and why that's the answer

The instinct is a geospatial index โ€” S2 cells, geohash, "so we can query by area." Resist it. Nobody queries "which devices are near this point"; the only query is "give me the reports for my device." That's a point lookup on a key the owner can compute. Since the sole access pattern is keyโ†’blob, you can afford the strongest possible privacy posture: the server holds ciphertext it cannot decrypt, keyed by a rotating identifier it cannot link to a user.

Store shape: a wide-column / KV store sharded by report_key, 7-day TTL, no secondary indexes at all. Writes are blind appends; reads are a multi-get of the ~2016 candidate keys the owner derives locally. The device registry and grants live in a normal replicated SQL store, because those do need joins, revocation semantics, and audit โ€” but they contain no location.

3. API interfaces

Report path (huge volume, unauthenticated by design)

POST /v1/reports        // from ANY device, incl. finders
{ report_key, ciphertext, ts }
โ†’ 204

// No auth header tied to identity โ€” attaching one
// would let the server correlate reporter to subject.
// Abuse control is anonymous rate-limiting:
// device attestation token + per-token quota.

POST /v1/commands/{device_id}
Signed-By: owner_key
{ type: LOCK, nonce, expires_at, signature }
โ†’ 202 { command_id }

Query path (tiny volume)

POST /v1/reports:fetch
{ report_keys: [k1..kN] }     // client-derived
โ†’ { found: [{key, ciphertext, received_at}] }
// client decrypts locally; server sees nothing

GET  /v1/devices               // my devices + status
POST /v1/grants  { grantee, device_id, expires_at }
DELETE /v1/grants/{id}
GET  /v1/audit?device_id=      // who looked, when

Fetch is a multi-get of opaque keys. The server can't tell whose device it is, or even that the N keys belong to one device. That property is the whole design.

4. Architecture

Lost device BLE beacon Finder device any stranger's Report ingest attestation check, anon rate limit Report store (KV) key=H(rot_pubkey) ciphertext, TTL 7d Registry (SQL) devices, grants, audit Command queue signed, nonce'd Fetch API multi-get by key no decrypt Push (FCM/APNs) wake device Owner client derives keys decrypts locally signs commands BLE authz + audit signed LOCK / WIPE

Green is the crowd-sourced report path โ€” a stranger's phone relays an encrypted blob it cannot read, to a server that also cannot read it. Red is the command path, signed by the owner's key and verified on the device, so a compromised server still can't wipe your phone.

5. Deep dive โ€” the privacy construction

Rotating keys, derived by both sides

At pairing (owner device + lost device share a secret):
  master_secret  ms          (never leaves either)

Every 15 minutes, epoch i:
  sk_i  = KDF(ms, i)                  # both sides derive
  pk_i  = curve_pub(sk_i)             # advertised over BLE
  key_i = SHA256(pk_i)                # the report_key

Finder:
  sees pk_i over BLE
  loc_ct = ECIES_encrypt(pk_i, {lat,lng,acc,ts})
  POST { report_key: SHA256(pk_i), ciphertext: loc_ct }

Owner looking for the device:
  for i in last 7 days (672 epochs @15m):
      keys.append(SHA256(pk_i))
  POST /reports:fetch { keys }        # ~672 keys
  decrypt each hit with sk_i          # locally

What each property buys, and what it costs

  • Server can't decrypt โ€” it never has any private key. A full database dump leaks nothing but "some blobs existed."
  • Server can't link reports to a device โ€” pk_i rotates every 15 min, so consecutive reports look unrelated. This is what stops the server (or an insider) reconstructing anyone's movement history.
  • Finder can't track the lost device โ€” it sees only a rotating pseudonym and can't decrypt what it just uploaded.
  • Stalking defense โ€” receiver devices detect an unknown beacon that persists across epochs while moving with you and alert the user. This is a real, shipped requirement and volunteering it is a strong signal.

Cost, stated honestly: ~672 key lookups per query instead of one; no server-side features at all (no "notify me when it moves", no geofencing, no server-side history search); and if the owner loses every device holding the master secret, past reports are unrecoverable. That last one is a genuine product trade and worth naming.

6. Deep dive โ€” command authenticity (the wipe problem)

Sequence

1. Owner authenticates (with 2FA) to the account service.
2. Client builds:
     cmd = {device_id, WIPE, nonce, expires_at}
     sig = sign(owner_private_key, cmd)
3. Server stores cmd, does NOT need to trust itself:
     it also checks the owner's session, but that is
     defense in depth, not the control.
4. Push wakes the device.
5. DEVICE verifies sig against the owner pubkey it
   pinned at enrollment, checks nonce not seen and
   expires_at in the future, then executes.
6. Device ACKs with a signed receipt -> audit log.

Why verification must happen on the device

If the server decides who may wipe, then a server compromise, a malicious insider, or a support-tool bug can wipe 10โธ devices. Verifying the owner's signature on the device reduces the server to an untrusted transport: the worst it can do is refuse to deliver commands (a denial of service, recoverable) rather than forge them (catastrophic, irreversible).

Defend it โ€” replay. Without nonce + expires_at, anyone who captures a valid LOCK command can replay it forever. The device keeps a small seen-nonce window and rejects anything expired, which bounds that memory.

Defend it โ€” the offline wipe race. A stolen device that's kept offline never gets the command. WIPE is queued with a long TTL and fires the moment it gets network; meanwhile the real protection is local full-disk encryption, which the remote wipe merely accelerates by destroying the key. Say this โ€” candidates who present remote wipe as the primary protection are wrong, and the interviewer knows it.

7. Deep dive โ€” efficiency (the third named axis)

Battery on the device

ChoiceWhy
BLE advertise, don't GPS-fixThe lost device only broadcasts a rotating pubkey โ€” no GPS, no radio uplink. Micro-amp scale; a dead-battery phone can beacon for hours on reserve power.
Finder supplies the locationThe finder was already computing its own position. Marginal cost โ‰ˆ one small HTTPS POST, batched.
Batch + opportunistic uploadFinders buffer reports and flush on Wi-Fi / when the radio is already up. Never wake the modem just for this.
Adaptive rate on the owner's own devicesReport position on significant-change, not on a timer. Stationary device โ†’ near-zero traffic.

Efficiency on the server

Write path is 10โถ/s of ~200-byte blind appends with no index maintenance and no read-modify-write โ€” the cheapest possible write. Sharded by report_key, which is a hash, so distribution is uniform by construction and there is no hot key possible: a popular location doesn't concentrate, because keys are per-device-epoch, not per-place.

Dedup: twenty phones in a train carriage all report the same beacon in the same epoch. Twenty near-identical ciphertexts under one key. Cap at ~K reports per (key, epoch) โ€” first-K-wins โ€” and drop the rest at ingest. That's a ~10ร— write reduction in dense areas and it costs nothing, since the owner only needs one good fix per epoch.

7-day TTL is doing heavy lifting: it bounds storage, and it bounds the blast radius of any future cryptographic weakness. Longer retention has no product value โ€” nobody finds a phone with a 3-week-old fix.

8. Follow-ups โ€” answers to have ready

A malicious actor floods the report endpoint with garbage under someone's key.

They'd have to know H(pk_i), which requires BLE proximity during that 15-minute epoch โ€” so the attack is inherently local and short-lived. Beyond that: the endpoint requires a device-attestation token (Play Integrity / DeviceCheck equivalent) with a per-token quota, so anonymous doesn't mean unlimited. And garbage ciphertext simply fails to decrypt on the owner's client and is discarded โ€” the poisoning doesn't corrupt anything, it just adds noise, which the K-per-epoch cap bounds.

The report store dies. What breaks?

Fail-open on the write path: ingest keeps accepting and buffers to a durable log, because a dropped report is permanently lost โ€” there's no retry from a stranger's phone that has already walked away. Reads simply return nothing, and the client shows "no recent location" rather than a stale one, which matters: showing an hour-old position as current sends someone to the wrong address. Registry and command paths are on a separate store, so lock and wipe still work during a location outage โ€” that separation is deliberate and worth calling out.

How does family sharing work if the server can't decrypt?

Key sharing, not server-side authorization. The owner wraps the master secret (or a derived, time-boxed viewing key) to the grantee's public key and hands it over through the registry. The server stores an opaque wrapped blob and enforces the grant record for UX and audit, but the actual cryptographic capability lives with the grantee. Revocation is therefore not instant for already-fetched data โ€” you must rotate the master secret on revoke, which invalidates future reports but not the ones the grantee already downloaded. State that limitation plainly; pretending revocation is instant is the wrong answer.

Law enforcement requests a user's location history.

The design's answer is that we can't produce it โ€” we hold ciphertext keyed by unlinkable rotating identifiers, with a 7-day TTL. That's not evasion, it's an architectural property chosen up front, and it's the same posture Apple ships. What we can produce is the registry: account, device enrollment, and the audit log of who queried. Being explicit about where the boundary sits is the mature answer here.

672 keys per fetch seems wasteful. Optimize it?

In practice the client fetches newest-first and stops early โ€” the common case is "where is it right now," which hits within the first few keys, so the tail only gets scanned when a device has been offline for days. You can also widen the epoch (15 min โ†’ 1 h) to cut keys 4ร—, at the cost of coarser unlinkability, since a longer-lived pseudonym is easier to correlate. That's the exact trade to name: epoch length is the privacy/efficiency dial. And the fetch is a batched multi-get on a hash-sharded KV store, so 672 keys is a handful of parallel shard reads, not 672 round trips.

Scale to 10ร— devices.

The write path scales linearly โ€” hash-sharded blind appends with no cross-shard coordination, so 10โท/s is more shards and nothing else. The parts that don't scale for free are the anonymous abuse quota (attestation token issuance becomes the bottleneck) and BLE spectrum congestion in dense areas, which is a client-side problem solved by advertise-interval backoff. The server design genuinely doesn't need re-architecting, and being able to say why โ€” no indexes, no joins, no hot keys, uniform hashing โ€” is the point.

9. Numbers to drop

Write volume

  • 10โถ reports/s ร— ~250 B = 250 MB/s ingest, ~21 TB/day.
  • With per-epoch dedup at K=5 in dense areas: ~10ร— reduction โ†’ ~2 TB/day effective.
  • 7-day TTL โ†’ ~15 TB steady state. Small enough to keep entirely on SSD.
  • Shard at ~50k writes/s/node โ†’ ~20 ingest shards at peak. Uniform by construction.

Read & crypto

  • Queries: say 10โธ users ร— 1 check/week โ‰ˆ 165 QPS. Three orders of magnitude below writes โ€” which justifies optimizing everything for write.
  • 672 keys/fetch ร— 165 QPS = ~110k key lookups/s. Trivial for a hash-sharded KV.
  • ECIES encrypt on a finder: ~1 ms, negligible against a BLE scan already running.
  • Key derivation on the owner client: 672 ร— KDF โ‰ˆ tens of ms, done once per open.

10. 30-second recap

The lost device does nothing but broadcast a Bluetooth pseudonym derived from a shared master secret, rotating every 15 minutes. Any nearby stranger's phone picks it up, encrypts its own GPS position to that rotating public key, and posts the blob under a key that is just the hash of the pseudonym. The server stores ciphertext it cannot decrypt, under identifiers it cannot link to a user or to each other, with a 7-day TTL โ€” so a full breach leaks nothing and there's no location history to subpoena or for an insider to abuse. The owner derives the same key sequence locally, multi-gets the last week of candidate keys, and decrypts on-device. Commands like lock and wipe are signed with the owner's key and verified on the device, so a compromised server can delay a wipe but never forge one. Efficiency falls out of the same design: blind appends with no index, hash-uniform sharding so no hot keys, and the beacon costs microamps because the finder supplies the location. The price I'm paying is real โ€” 672 key lookups instead of one, no server-side geofencing or history, and revocation that requires rotating the master secret rather than flipping a flag.

Street View image ingest full

L6/L7, Google Cloud Storage org, phone screen, 2026-01. Verbatim: "Design a system that supports Google Map Street View storage. The images will be uploaded from a taxi. Each taxi has a camera; images are processed by downstream โ€” image understanding, display to users, map generation." Reported follow-ups: auth design and what to do when a token is compromised; upload response design; behavior on poor network. Interviewer was near-silent and let the candidate drive for the full hour.

1. Requirements in one line

Functional

  • Fleet vehicles capture panoramic frames + GPS/IMU metadata and upload them.
  • Uploads survive intermittent connectivity: resumable, deduplicated, eventually complete.
  • Ingested images are durably stored and published to downstream consumers (blur/PII, image understanding, tiling, map generation).
  • Operators can see per-vehicle ingest status and re-drive coverage gaps.

Out of scope: the ML models themselves, the map-tile serving CDN, camera firmware.

Non-functional

  • 10k vehicles; each ~1 frame/sec while driving, ~8 h/day. ~10 MB/panorama.
  • Ingest is throughput-bound, not latency-bound โ€” nobody needs the photo in a second.
  • Durability is the hard requirement: a lost frame means re-driving a street. Target 11 nines.
  • Cost matters at petabyte scale; storage tiering is part of the design, not an afterthought.

Core tension: the client is a moving vehicle on flaky cellular, generating far more data than its link can carry, and it cannot be trusted (a stolen credential lets someone poison the map). So the design is really about a well-behaved offline-first client plus an untrusted-upload security model, with the storage tier being the easy part.

2. Core entities & schema

CaptureSession
  session_id, vehicle_id, driver_id, started_at, route_id
  status -- ACTIVE | UPLOADING | COMPLETE | PARTIAL

Frame                       -- metadata, in a wide-column store
  frame_id = uuidv7(capture_ts)      -- time-sortable
  session_id, seq
  content_hash (sha256 of the image bytes)   <-- dedup + integrity
  captured_at, lat, lng, heading, accuracy
  s2_cell_l16                                 -- geo index key
  blob_uri, bytes, camera_id
  state -- PENDING | UPLOADED | VERIFIED | PUBLISHED | QUARANTINED

Upload                      -- the resumable session
  upload_id, frame_id, total_bytes
  received_ranges []        -- byte ranges committed
  expires_at

Blob                        -- object storage, immutable
  path = /raw/{s2_cell}/{yyyymmdd}/{content_hash}
  // content-addressed: same bytes uploaded twice = one object

Why object storage + a metadata store, and why content-addressing

Images never belong in a database โ€” 10 MB blobs are what object storage exists for: cheap, replicated, tiered, with no query needs. Metadata is small, queried by geography and by session, and needs a secondary index, so it lives separately in a wide-column store keyed by frame_id with an s2_cell index for "what have we covered here?"

Content-addressing the blob path is the load-bearing choice. The client hashes bytes before upload; the path is the hash. That gives you three things for free: (1) a retry that re-uploads the same frame overwrites itself instead of duplicating โ€” idempotency without any server-side dedup logic; (2) end-to-end integrity, since the server recomputes the hash and rejects a mismatch, catching both corruption and tampering; (3) natural dedup when a vehicle stops at a light and the camera captures near-identical frames โ€” though see the follow-ups, byte-identical is rarer than you'd hope.

uuidv7 for frame_id rather than random UUIDv4 because it's time-prefixed and therefore time-sortable, which makes session-ordered scans a contiguous range read rather than a scatter.

3. API interfaces

Upload path (resumable, the actual question)

POST /v1/uploads:init
Authorization: Bearer {short-lived vehicle token}
{ session_id, seq, content_hash, bytes, captured_at,
  lat, lng, heading }
โ†’ 200 { upload_id,
        upload_url,        // signed, direct-to-storage
        expires_in: 3600,
        already_have: false }   // <-- dedup short-circuit

PUT {upload_url}
Content-Range: bytes 5242880-10485759/10485760
โ†’ 308 { received: "0-10485759" }   // resume point
โ†’ 201 when complete

GET /v1/uploads/{upload_id}   // "where was I?"
โ†’ { received_ranges: [[0, 5242879]] }

Downstream & ops

// consumers subscribe, they don't poll
topic: frames.uploaded  { frame_id, blob_uri, geo }
topic: frames.published { frame_id, after blur+verify }

GET /v1/coverage?s2_cell=&since=
โ†’ { frames, gaps: [{lat,lng,reason}] }
GET /v1/sessions/{id}/status
โ†’ { captured: 28800, uploaded: 28112, pending: 688 }
POST /v1/frames/{id}:quarantine  { reason }

Answering the "upload response design" follow-up: return 308 with the exact committed byte range, plus already_have on init so a client that crashed after uploading never re-sends 10 MB. The response's job is to make the client's retry decision unambiguous.

4. Architecture

Vehicle agent local disk spool hash + queue Wi-Fi preferred Ingest control authz, init, resume Signed URL direct to bucket Object store (raw) content-addressed regional, versioned Frame metadata wide-column + S2 idx Verifier rehash, sanity-check Event log frames.uploaded Blur / PII pipeline faces, plates โ€” gate Image understanding signs, storefronts, lanes Tiling / pano stitch โ†’ user-facing CDN Map generation geometry extraction Lifecycle: hot 30d โ†’ nearline 1y โ†’ coldline archive raw kept forever (reprocessable); derived rebuilt on demand init PUT bytes

Bytes never traverse the ingest service โ€” it only issues signed URLs, so the control plane stays small and cheap while 10 MB objects go straight to storage. Downstream consumers subscribe to an event log rather than being called synchronously, so a slow ML pipeline can never back-pressure a vehicle on a cellular link. Note the blur/PII stage is a gate, not a parallel consumer.

5. Deep dive โ€” the offline-first client (the "poor network" follow-up)

The client is the interesting system

Capture loop (never blocks on network):
  frame -> local SSD spool + WAL entry
  compute sha256 while writing
  spool is a bounded ring: 500 GB โ‰ˆ 14 h of capture

Upload loop (independent, opportunistic):
  while spool not empty:
    pick oldest UNSENT frame
    if link_quality == CELLULAR and not urgent:
        upload at throttled rate (leave headroom)
    if link == WIFI (depot / known AP):
        upload at full rate, drain aggressively
    on 5xx / timeout:
        exponential backoff + jitter, keep frame
    on 201:
        mark SENT, delete local copy lazily
        (keep until VERIFIED event or 24h)

Chunk size adapts: start 8 MB, halve on each
timeout down to 256 KB, grow back on success.

Why this shape

Capture and upload must be decoupled queues. If capture blocks on upload, a tunnel or a dead cell means lost coverage and a re-drive โ€” the single most expensive failure in this system. Spooling to local disk turns a network problem into a storage problem, which is far cheaper.

Adaptive chunk size is the concrete answer to "network is bad." On a marginal link, a 10 MB PUT that dies at 90% wastes 9 MB of the vehicle's uplink; 256 KB chunks lose almost nothing per failure. The trade is more round-trips and more per-request overhead, so you grow the chunk back when the link proves itself.

Deletion is delayed until VERIFIED, not until the 201. A 201 means bytes landed; it doesn't mean the server rehashed them successfully. Deleting the only copy on a 201 is how you lose data to a silent corruption. Cost: local disk holds a day of already-uploaded frames.

Belt-and-suspenders โ€” the bounded spool. 14 hours of buffer exceeds an 8-hour shift, so a vehicle that never sees network all day still loses nothing; it drains at the depot on Wi-Fi. If the ring does wrap, drop by lowest marginal coverage value (a frame from an already-well-covered cell) rather than oldest-first. Volunteering a smarter eviction policy than FIFO reads well.

6. Deep dive โ€” auth, and the compromised-token follow-up

Layered credentials

Layer 1 โ€” hardware identity (long-lived, never sent)
  Per-vehicle key in a TPM/secure element at fleet
  provisioning. Private key cannot be extracted.

Layer 2 โ€” short-lived access token (the thing on the wire)
  Vehicle signs a challenge with the TPM key
    -> token service returns JWT, TTL 1 hour,
       scoped: {vehicle_id, session_id, WRITE only,
                geo_bbox of today's route}
  Never a long-lived API key on the device.

Layer 3 โ€” per-object signed URL (TTL 1 hour, single path)
  Scoped to exactly /raw/{cell}/{date}/{content_hash}
  So even a leaked URL can write ONE object,
  whose path is the hash of its own contents โ€”
  so it cannot write anything but that exact image.

"The token is compromised. What do you do?"

Answer in the order the interviewer is grading: contain, revoke, detect, prevent.

  • Blast radius is already small by design โ€” the token is write-only, expires in an hour, is scoped to one vehicle and one geographic box, and the signed URLs it can obtain are content-addressed. An attacker cannot read anyone's imagery, cannot overwrite existing frames with different bytes, and cannot write outside today's route area.
  • Revoke: add the token's jti to a deny list checked at the ingest control plane (small, because TTL is short โ€” the list self-cleans in an hour). Revoke the vehicle's provisioning cert to stop it minting new tokens.
  • Detect: the signals that catch this are behavioral โ€” uploads from an IP/ASN that doesn't match the vehicle's cellular carrier, frames whose GPS is inconsistent with the session's route continuity or physically impossible speed, upload rate exceeding the camera's capture rate, or two concurrent sessions for one vehicle_id. Any of these quarantines the frames rather than dropping them, so a false positive is recoverable.
  • Prevent: mutual TLS with the device cert on the control plane, so the bearer token alone is insufficient; anything anomalous goes to QUARANTINED and is human-reviewed before publish.

The point to land: the right answer to a compromised credential isn't a faster revocation loop, it's designing so that a compromised credential can't do much. Short TTL, narrow scope, content-addressing, and write-only are all doing that work.

7. Deep dive โ€” downstream fan-out and the PII gate

Event-driven, not synchronous

Ingest publishes frames.uploaded and stops caring. Consumers โ€” blur, image understanding, tiling, map generation โ€” subscribe independently with their own offsets, their own scaling, and their own failure domains. A stalled ML pipeline creates consumer lag, not upload failures.

Reprocessing is the reason raw is immutable and kept forever. Models improve; when a better sign-detection model ships you replay the last N years of frames.uploaded from the log (or scan the metadata store by S2 cell) rather than re-driving the planet. Derived artifacts โ€” tiles, extracted geometry โ€” are therefore disposable, and treating them that way is what makes the storage bill defensible.

The blur stage must be a gate, not a peer

Faces and license plates must be blurred before anything user-facing exists. If blur is just another parallel subscriber, there's a window where the tiling pipeline has already published an unblurred panorama โ€” a privacy incident with legal consequences in the EU and elsewhere.

So: two topics. frames.uploaded is consumed only by blur/verify. Blur writes a redacted derivative and emits frames.published; everything user-facing subscribes to that one. Internal-only consumers may read raw under separate authorization and audit.

Defend it: blur models miss things, so you need a takedown path โ€” a user reports their face, you re-blur that frame and purge derived tiles, which is only tractable because frame_id โ†’ tile lineage is recorded. Without lineage tracking, a single takedown request means rebuilding a region. Say this; it's the kind of operational detail that separates L5 from L4 here.

8. Follow-ups โ€” answers to have ready

Two frames have the same content hash. Is that dedup or a bug?

Usually a bug worth investigating. Real-world sensor noise means two genuinely separate captures essentially never produce identical bytes โ€” so a collision means either the camera emitted a duplicate frame (stuck buffer), or the client retried and re-uploaded, or someone is replaying captured data. The write is idempotent either way, so nothing breaks; but I'd count it as a metric and alert on the rate. For perceptual dedup of a vehicle stopped at a red light, content hashing is the wrong tool โ€” use a perceptual hash or simply drop frames when GPS shows the vehicle hasn't moved, which is cheaper and done on the client.

Object storage in one region goes down mid-shift.

Vehicles keep capturing โ€” that's the whole point of the local spool โ€” and uploads fail into backoff. The ingest control plane should detect the regional failure and start issuing signed URLs for a secondary region, since blobs are content-addressed and the region is just a prefix; the metadata store records where each blob actually landed. Recovery cost is an asynchronous cross-region copy job later. What you must not do is let the vehicle agent block or drop frames while the server sorts itself out.

A vehicle uploads frames with plausible but wrong GPS.

This is the poisoning case and it's the reason for the verifier stage. Checks: route continuity (does this position follow the previous one at a physically possible speed?), IMU cross-check, and consistency with existing imagery for that S2 cell. Suspicious frames go to QUARANTINED โ€” never auto-deleted, because a false positive on a legitimately unusual route (a ferry, a new road) would silently lose coverage. Quarantine plus human review is the right severity.

How much does this cost, and how do you cut it?

Raw is the dominant line item and it grows forever. Lifecycle tiering does the heavy lifting: hot for 30 days while the pipelines chew on it, nearline for a year, coldline archive after. Raw is never deleted because re-driving costs vastly more than storing. Derived artifacts go the other way โ€” regenerate tiles on demand from raw rather than storing every historical version. And compression: raw sensor frames stored as lossless for reprocessing, but the panorama derivatives as aggressive lossy, since human viewers can't tell.

Why not have the vehicle upload straight to the bucket with no control plane at all?

Because you'd need a long-lived broad credential on a physically accessible device โ€” exactly the thing the layered auth avoids. The control plane exists to mint narrow, short-lived, per-object grants and to record metadata transactionally with the intent to upload, so you can detect a frame that was initiated but never completed. It stays cheap because it never touches bytes: 10k vehicles at ~1 init/sec is ~10k QPS of tiny JSON, which is a handful of pods.

Scale from 10k vehicles to 1 M (consumer dashcams).

The storage and event tiers scale horizontally without change. Three things break. First, trust: consumer devices have no TPM and no fleet provisioning, so poisoning becomes the dominant risk and you need reputation scoring and heavy cross-validation between contributors. Second, redundancy: a million dashcams will cover the same popular streets thousands of times, so you need server-side coverage-aware admission โ€” reject the upload at init time when that cell is already well covered, which also saves the contributor's bandwidth. Third, cost per useful frame collapses unless you do that admission control. That's the re-architecture point, and it's on the ingest control plane, not on storage.

9. Numbers to drop

Ingest volume

  • 10k vehicles ร— 1 frame/s ร— 8 h = 288 M frames/day.
  • At 10 MB each: ~2.9 PB/day raw. This number is the design constraint โ€” say it early.
  • Sustained ingest if spread over 24 h: ~33 GB/s. Per vehicle that's only ~10 Mbps, which is exactly why depot Wi-Fi drain matters โ€” cellular can't carry it all.
  • Control plane: ~3.3k init QPS average, ~10k peak. Trivially small next to the byte volume.

Storage & metadata

  • ~1 EB/year raw. With lifecycle tiering (30d hot / 1y nearline / then coldline) the blended rate is roughly 4โ€“5ร— cheaper than all-hot.
  • Metadata: 288 M rows/day ร— ~500 B = ~145 GB/day โ€” 0.005% of the blob volume, which is why it can afford indexes.
  • Vehicle spool: 500 GB SSD รท 10 MB = 50k frames โ‰ˆ 14 h of capture, comfortably more than one shift.
  • Depot drain: 14 h of backlog at 1 Gbps Wi-Fi โ‰ˆ 70 min. Sizes the depot link.

10. 30-second recap

The vehicle agent treats capture and upload as two independent queues joined by a local disk spool sized for a full shift, so a tunnel or a dead cell never costs coverage โ€” that's the expensive failure here, since a lost frame means re-driving a street. The client hashes each frame, calls a small control plane to initialize a resumable upload, and gets back a signed URL scoped to exactly one content-addressed path; bytes then go straight to object storage, so the control plane never touches a 10 MB payload. Content-addressing gives idempotent retries, end-to-end integrity, and free dedup. On a bad link the client halves its chunk size down to 256 KB so a failure wastes almost nothing, and it doesn't delete its local copy until the server confirms it rehashed the bytes. Auth is layered โ€” a TPM key that never leaves the vehicle mints one-hour, write-only, geo-scoped tokens โ€” so a stolen token can't read anyone's imagery, can't overwrite existing frames, and expires before you've finished revoking it. Ingest publishes an event; blur and PII redaction is a gate rather than a parallel consumer, so nothing user-facing can ever see an unblurred face, and everything downstream subscribes to the post-blur topic. Raw is immutable and kept forever because models improve and reprocessing is far cheaper than re-driving.

Global chain restaurant menu sync full

็‹—ๅฎถ SD, 2026-02. Verbatim: "็ป™ไธ€ไธชๅ…จ็ƒ่ฟž้”้คๅŽ…่ฎพ่ฎก่œๅ•ๆ›ดๆ–ฐ็ณป็ปŸใ€‚ๅ‡่ฎพๆฏๅฎถ้คๅŽ…ๅฏ่ƒฝๆœ‰ๅ„็งๅฑ•็คบ่œๅ•็š„่ฎพๅค‡ใ€‚ ่ฟž้”้คๅŽ…ๅœจๅ„ๅ›ฝๅฎถ็š„่œๅ•ๆœ‰็›ธๅŒ็š„้ƒจๅˆ†๏ผŒไนŸๆœ‰ๆœฌๅœฐๅŒ–็š„้ƒจๅˆ†ใ€‚่œๅ•็š„ๆ›ดๆ–ฐ็”ฑๆ€ป้ƒจๆŽงๅˆถ๏ผŒๆฏๅคฉๆ—ฉไธญๆ™š้คๅฏ่ƒฝ่ทŸๅ‰ไธ€ๅคฉไธไธ€ๆ ทใ€‚" Read the prompt carefully โ€” it hands you three hard requirements: heterogeneous display devices, inheritance with localization, and time-scheduled variants.

1. Requirements in one line

Functional

  • HQ authors a global menu; countries/regions/stores override parts of it (price, availability, language, local items).
  • Menus vary by daypart: breakfast / lunch / dinner, and can differ day to day.
  • Push updates to heterogeneous devices: digital menu boards, kiosks, tablets, drive-thru displays, the mobile app, and third-party delivery aggregators.
  • Scheduled activation ("this menu goes live Monday 6am local") and emergency rollback.
  • A store must show a correct menu even with no network.

Out of scope: order placement, inventory/POS integration beyond an availability flag, the authoring UI itself.

Non-functional

  • 40k stores ร— ~5 devices = 200k endpoints, ~100 countries.
  • Update propagation: minutes is fine for planned changes; seconds matter for a takedown (allergen recall, wrong price).
  • Availability >> consistency at the edge: a device must never show a blank screen.
  • Correctness of price is the one thing that's legally and financially sensitive.

Core tension: the update volume is laughably small โ€” this is not a scale problem. It's a configuration-correctness and safe-rollout problem across 200k unreliable, heterogeneous, intermittently-connected endpoints. Candidates who reach for Kafka and sharding here have misread the question.

2. Core entities & schema

MenuItem                  -- the catalog (global)
  item_id, canonical_name, category, image_set,
  allergens[], nutrition{}, default_price_cents

MenuLayer                 -- the inheritance chain. THE key entity.
  layer_id, scope_type -- GLOBAL | COUNTRY | REGION | STORE
  scope_id             -- null | "JP" | "kanto" | "store_8842"
  parent_layer_id
  patch {              -- sparse: only what this layer CHANGES
    add:      [item_id...],
    remove:   [item_id...],
    override: { item_id: {price, name_i18n, image} }
  }

MenuVersion               -- immutable, content-addressed
  version_id = hash(resolved menu bytes)
  layer_set[], daypart, effective_from, effective_to
  state -- DRAFT | APPROVED | STAGED | ACTIVE | ROLLED_BACK
  approved_by, approved_at

DeviceRegistration
  device_id, store_id, type (BOARD|KIOSK|DRIVE_THRU|APP),
  capabilities {screen, aspect, supports_video, locales},
  current_version_id, last_heartbeat_at

Why layered patches instead of a menu per store

The naรฏve model โ€” one full menu document per store โ€” means 40k documents, and a global price change requires rewriting all of them with no atomicity and no way to tell an intentional local difference from a stale copy. The layered model makes inheritance explicit and diffable: HQ edits the GLOBAL layer, and every store that hasn't overridden that item inherits the change automatically.

Resolution is a deterministic fold โ€” GLOBAL โ†’ COUNTRY โ†’ REGION โ†’ STORE, later layers win. Determinism is what lets you content-address the resolved output: same layers in, same version_id out. That in turn gives you free change detection (a device compares hashes), free dedup across identical stores, and an unambiguous rollback target.

Storage is boring on purpose: a replicated SQL store for layers and versions (you need transactions, approval workflow, and audit on price), plus object storage + CDN for the resolved bundles that devices actually download. Resolution happens once at publish time, not per-device โ€” 40k stores ร— 3 dayparts is ~120k resolutions per publish, which is a few seconds of compute, and it means a device does zero logic.

3. API interfaces

Device path (pull, not push)

GET /v1/devices/{device_id}/manifest
If-None-Match: {current_version_id}
โ†’ 304 Not Modified                      // the 99.9% case
โ†’ 200 {
    schedule: [
      { daypart: "breakfast", from: "06:00",
        to: "10:30", version_id: "a1b2",
        bundle_url: "https://cdn/.../a1b2.tar",
        sha256: "...", size: 4_200_000 },
      { daypart: "lunch", ... }
    ],
    poll_after_sec: 300,
    emergency_channel: "wss://..."
  }

POST /v1/devices/{device_id}/heartbeat
{ active_version_id, staged_version_ids[],
  render_errors[], clock_skew_ms }

Devices pull on a jittered interval and hold a long-lived websocket only for emergencies. Pull scales trivially through a CDN and survives a device being offline for a week; push-only would require the server to track 200k connections just to deliver a change that isn't urgent.

Authoring & ops path

PUT  /v1/layers/{layer_id}      { patch }
POST /v1/versions:preview
  { layer_set, daypart, store_id }
  โ†’ fully resolved menu + diff vs current

POST /v1/versions/{id}:approve   // 2-person for price
POST /v1/rollouts
  { version_id, cohort: "canary_50_stores",
    effective_from: "2026-03-01T06:00 LOCAL" }
POST /v1/rollouts/{id}:halt
POST /v1/emergency/takedown
  { item_id, scope, reason }     // seconds, not minutes

GET  /v1/fleet/status?version_id=
  โ†’ { on_target: 39_812, stale: 188, unreachable: 40 }

4. Architecture

HQ authoring + country / store Layer store (SQL) patches, approvals, audit on price Resolver fold layers โ†’ content-hash bundle Validator price/allergen gates Bundle store immutable objects per version_id Rollout controller cohorts, halt, rollback CDN edge-cached bundles + manifest Emergency ws takedown push Heartbeat / fleet who is stale? Store gateway local cache of ALL staged versions + local clock Devices board / kiosk / drive-thru / app render from capability only if valid pull

Publishing is a fold + validate + freeze into an immutable, content-addressed bundle. Distribution is boring CDN pull with a per-store gateway that pre-stages every upcoming version, so activation is a local event that needs no network at all. The red path is the only push channel, reserved for takedowns.

5. Deep dive โ€” staging and local activation (the daypart problem)

Never activate over the network

Publish (T-24h or more):
  resolver produces bundles for every
  (store, daypart) combination for the next 48h
  โ†’ manifest lists them ALL with local times

Store gateway (continuously):
  downloads every staged bundle
  verifies sha256
  keeps last-known-good + all staged

Activation (exactly at 06:00 store-local):
  gateway compares local clock to schedule
  swaps active pointer atomically
  devices re-render from local files

Network down at 06:00?  Still switches.
Network down for 3 days? Runs out of staged
  versions โ†’ holds last-known-good, alarms.

Why this is the crux of the question

The prompt's "ๆ—ฉไธญๆ™š้คๅฏ่ƒฝ่ทŸๅ‰ไธ€ๅคฉไธไธ€ๆ ท" is a trap for anyone who designs a push-at-activation-time system. If the breakfast menu is pushed at 06:00 local, then every store in a timezone hits you simultaneously โ€” a thundering herd at every hour boundary around the globe โ€” and any store whose link is down at that moment shows the wrong menu during its busiest period. Pre-staging converts a synchronized network event into an unsynchronized download plus a local clock comparison.

Defend it โ€” clock skew. The design now depends on device clocks. A gateway with a drifting clock switches to dinner at lunchtime. Mitigations: NTP on the gateway, clock skew reported in every heartbeat with an alert above ~60 s, and a sanity rule that the gateway won't jump more than one daypart forward without a server confirmation. This is the failure mode the interviewer will hunt for once you propose local activation, so raise it yourself.

Timezones and DST: schedules are stored as local wall-clock times with an IANA timezone ID per store, never as UTC offsets โ€” offsets break twice a year, and a store in a DST-shifting region would show breakfast an hour late every spring.

6. Deep dive โ€” heterogeneous devices

Ship data, plus per-capability renditions

Bundle contents:
  menu.json          -- semantic content, device-agnostic
                        items, prices, i18n strings, order
  layouts/
    board_16x9.json     kiosk_portrait.json
    drive_thru.json     app_compact.json
  assets/
    {hash}.webp @1x/2x   {hash}.mp4 (boards only)
  render_hints: { max_items_per_panel, font_scale }

Device fetches ONLY the layout + assets matching
its registered capabilities โ†’ a drive-thru display
never downloads 4K video it can't show.

The trade

The alternative โ€” pre-rendering images per device type server-side โ€” makes devices trivially dumb but multiplies bundle count by device type, and any layout tweak means republishing everything. Shipping semantic JSON plus layout descriptors means one content pipeline and per-device presentation, at the cost of some rendering logic on the endpoint.

Defend it โ€” the capability lie. Devices misreport capabilities, and old firmware won't understand a new layout field. So: version the bundle schema, require devices to declare a supported schema range in the manifest request, and have the server serve the highest bundle version that device can parse. Unknown fields must be ignored, not fatal. Without this, one new field bricks every unpatched screen in the fleet โ€” and menu boards get firmware updates roughly never.

Third-party aggregators (delivery apps) are just another consumer type, but they pull via API rather than bundle, and they're the one consumer where a stale price becomes a customer-facing refund. Give them a webhook on version change plus a short cache TTL.

7. Deep dive โ€” safe rollout and the takedown path

Two speeds, deliberately

Planned changeEmergency takedown
TriggerApproved version, scheduledAllergen error, wrong price, recall
PathCDN pull, jittered, stagedWebsocket push + poll interval collapses to 10 s
GranularityWhole versionSingle item_id: hide it
TargetMinutesโ€“hours< 30 s to 99% of fleet
Fallbackโ€”If push fails, next poll catches it

A takedown is a subtractive overlay, not a new version โ€” "hide item X everywhere" is a tiny message that any device can apply to whatever bundle it's currently running, including a stale one. That's why it's fast and why it works on devices that haven't synced.

Canary the way you'd canary code

A bad menu version is a production incident with revenue impact. Roll out by cohort: 50 stores โ†’ one region โ†’ country โ†’ global, with automatic halt on error signals (render failures in heartbeats, device version-adoption stalling, or a spike in POS price mismatches). POST /rollouts/{id}:halt freezes advancement; rollback re-points to the previous version_id, which is instant because bundles are immutable and still cached at the edge.

Validation gates before any of that: no item priced at 0 or negative; no price change greater than X% without a second approver; every item has allergen data in every locale it's published to; every referenced asset exists; total items fit the smallest target screen. These catch the realistic disaster โ€” a fat-fingered price rolled to 40k stores โ€” far more reliably than a canary does, because a canary only catches what the canary stores happen to sell.

Defend it โ€” partial fleet state. During any rollout, different stores are legitimately on different versions. That's fine for menus but not for a nationally advertised promotion, so the schedule's effective_from is what gates display, while the rollout gates distribution. Separating those two is what lets you pre-stage for days and still have every store flip together.

8. Follow-ups โ€” answers to have ready

A store's internet has been down for two days. What's on the screen?

The correct menu, because the gateway pre-staged 48 hours of bundles and activation is a local clock comparison. Past that horizon it holds last-known-good and keeps switching dayparts using the most recent schedule it has โ€” a slightly stale menu beats a blank board. The heartbeat gap raises the store in the fleet dashboard as unreachable, and there's a manual USB/local-admin path for a store manager to sideload a bundle in a true emergency. The design decision to state plainly: we choose availability over freshness at the edge, because a dark menu board stops sales entirely.

HQ changes a global price but Japan has overridden that item. What happens?

Japan keeps its override โ€” that's the whole point of the layer model, and it's the correct default. The risk is silent divergence: HQ thinks it changed a price globally and doesn't realize 12 countries have pinned it. So the authoring UI must show, at edit time, how many downstream layers override this item, and the publish diff must list them explicitly. Optionally offer a "force, clearing overrides" action, which requires a higher approval level. The technical model is easy here; the product affordance around it is what makes it correct.

The resolver has a bug and produces a garbage menu.

Validation gates run after resolution and before the bundle is frozen, so structurally-invalid output never gets a version_id. For semantically-wrong-but-valid output, the canary cohort plus adoption monitoring catches it, and rollback is instant because the previous immutable bundle is still in the CDN and in every gateway's local cache. Worth noting the resolver is a pure function of the layer set, so it's the easiest component in the system to test exhaustively โ€” golden-file tests over real store configurations.

Why not just push everything over websockets? 200k connections isn't that many.

It isn't, and you do need the channel for takedowns. But making it the primary path means the server owns delivery state for 200k endpoints, reconnect storms after any deploy, and no natural way to serve a device that was offline for a week. Pull with a CDN inverts that: the server publishes an immutable artifact and forgets, and the edge absorbs the load. Push is reserved for the case where seconds matter, which is rare, small, and subtractive.

How do you know the fleet is actually correct right now?

Heartbeats carry active_version_id, so the fleet endpoint answers "how many stores are on the version I intended, right now" โ€” and that ratio, not the publish succeeding, is the definition of a successful rollout. Alert on stale-device count and on adoption curves that plateau. Add a synthetic: a handful of reference devices that screenshot and diff against the expected render, which catches the class of bug where the data is right and the rendering is wrong.

Scale it to a franchise model where each owner can edit their own menu.

The layer model already supports it โ€” franchisees author at the STORE layer. What changes is governance: you now need per-layer permissions, a policy engine restricting which fields a franchisee may override (price yes, allergen data absolutely not, brand assets no), and much stronger validation because the authors are no longer a small trusted team. Authoring volume goes from a handful of HQ edits a day to tens of thousands, which finally makes the write path non-trivial โ€” but distribution is unchanged.

9. Numbers to drop

Fleet & distribution

  • 40k stores ร— ~5 devices = 200k endpoints; ~5 dayparts ร— 2 days staged = ~10 bundles live per store.
  • Manifest poll every 5 min, jittered โ†’ 200k / 300 s โ‰ˆ 670 QPS, and ~99.9% are 304s. This is a small service.
  • Bundle ~5 MB (mostly images, incremental after the first). Full fleet refresh = 200k ร— 5 MB = 1 TB, spread over hours via CDN. Trivial.
  • Resolution cost: 40k stores ร— 3 dayparts = 120k folds per publish, ~ms each โ†’ a few seconds on one machine.

Latency targets

  • Planned change โ†’ 99% of fleet: ~15 min (3 poll intervals).
  • Emergency takedown โ†’ 99%: < 30 s via websocket, with the 10 s collapsed poll as backstop.
  • Daypart switch: 0 ms of network, purely local.
  • Rollback: instant (re-point to a cached immutable bundle).
  • Data volume overall is so small that the entire design is driven by correctness and offline behavior โ€” say this explicitly, it shows you sized the problem.

10. 30-second recap

Menus are modelled as a chain of sparse patch layers โ€” global, country, region, store โ€” that resolve by a deterministic fold, so HQ edits propagate automatically to everyone who hasn't overridden that item, and local customization is a first-class concept rather than a copy. Resolution happens once at publish time and produces an immutable, content-addressed bundle per store and daypart, which passes validation gates for price sanity and allergen completeness before it's allowed a version ID. Distribution is CDN pull with jittered polling โ€” mostly 304s โ€” and each store gateway pre-stages 48 hours of upcoming bundles, so the breakfast-to-lunch switch is a local clock comparison that needs no network at all. That's the key move: it avoids a global thundering herd at every hour boundary and means a store with a dead link still shows the right menu. Devices get semantic JSON plus a layout matching their declared capabilities, with schema versioning so a new field doesn't brick unpatched menu boards. Rollout is canaried by cohort with automatic halt, and there's a separate fast path for emergency takedowns โ€” a subtractive "hide this item" overlay that applies to whatever bundle a device is already running, so it works in under 30 seconds even on a stale device. The thing I'd watch hardest is clock skew on the gateways, since local activation makes correctness depend on their clocks.

Dictionary range query store condensed

็‹—ๅฎถ L6, 2026-02. Verbatim: "่ฎพ่ฎกไธ€็งๅญ˜ๅ‚จ็ณป็ปŸ๏ผŒ่ƒฝๅคŸๆŒ‰ๅญ—ๅ…ธๅบๆŸฅ่ฏขไธ€ไธชๅŒบ้—ดๅ†…็š„ๆ‰€ๆœ‰่ฏใ€‚ ๆฏ”ๅฆ‚ๆ•ฐๆฎๆœ‰ {aa, aaa, ac, f, z}: ่พ“ๅ…ฅ [a, b] โ†’ {aa, aaa, ac}; ่พ“ๅ…ฅ [a, aa] โ†’ {aa}; ่พ“ๅ…ฅ [] โ†’ {f, z}; ่พ“ๅ…ฅ [ab, b] โ†’ {ac}." This is the rare SD prompt that's really a data-structure choice with a distributed-systems tail. Nail the local structure first, then scale it.

1. Requirements & the semantics trap

Clarify before designing

  • Boundary semantics: the examples imply [a, b) half-open โ€” [a, aa] returns {aa} but not {aaa, ac}, and [ab, b] returns {ac}. Confirm this out loud; getting it wrong invalidates everything after.
  • Result size: could a range match 10โธ words? If yes, the API must paginate with a cursor, not return a list.
  • Mutability: read-only corpus (build once, serve forever) or live insert/delete? This is the biggest fork in the design.
  • Scale: 10โถ words fits in memory on one box; 10ยนยน words does not. Ask.
  • Latency SLO, and whether prefix/autocomplete queries are also needed (they change the structure choice).

Assume, and say so

  • 10ยนโฐ words, ~20 bytes each โ†’ ~200 GB. Doesn't fit one machine.
  • Mostly reads; writes are bulk-loaded plus a live trickle.
  • p99 < 50 ms for a range returning โ‰ค 1000 results.
  • Results must be returned in sorted order and paginated.

Core tension: range queries want data sorted and contiguous; distribution wants data hashed and uniform. Those are directly opposed, and choosing sorted (range-partitioned) means accepting hot shards. Everything below follows from that.

2. Structure choice โ€” say why, not just what

OptionRange queryCostVerdict
Sorted array + binary searchO(log n) to find start, then scanImmutable; insert is O(n)Right for a static corpus. Say this first โ€” it's the simplest thing that works.
B+ treeO(log n) descend, then walk the leaf linked listNode splits on write; ~1.3ร— spaceThe answer for a mutable corpus. Leaves are sorted and linked, which is exactly a range scan.
LSM tree (SSTables)Merge-iterate across sorted runsRead amplification; compactionRight when writes dominate. Range scan touches every level.
Trie / radix treePrefix-natural; range needs bounded DFSHigh pointer overhead, poor cache localityOnly if prefix queries are the real requirement. For arbitrary [a, b) it's worse than a B+ tree.
Hash indexImpossible โ€” hashing destroys orderโ€”Naming why this fails is a cheap point. Say it.

Pick B+ tree for the mutable case and be explicit about why the trie loses: a trie shines when you're matching a shared prefix, but [ab, b) has no single prefix โ€” you'd descend to ab, DFS its subtree, then walk siblings up to b, which is a lot of pointer chasing for what a B+ tree does as one sequential leaf walk.

3. Architecture

Client range(a, b) Coordinator routing table โ†’ overlapping shards k-way merge Shard [ , c) B+ tree + replicas Shard [c, m) B+ tree + replicas Shard [m, ) B+ tree + replicas Range registry split / merge boundaries Rebalancer splits hot ranges Per-shard layout internal nodes in RAM leaves on SSD, linked + prefix- compressed skipped

Range-partition, don't hash-partition. The routing table maps lexicographic boundaries to shards, so range(a, b) touches only the shards that overlap [a, b) โ€” often one. Hashing would fan every query to every shard.

4. Trade-offs worth stating

The hot-shard problem you created

Real word distributions are wildly skewed โ€” a shard covering [s, t) holds far more English words than one covering [x, z), and query traffic is skewed too. Range partitioning guarantees this; it's the price of ordered scans.

  • Split by load, not by key space. Choose boundaries so each shard holds roughly equal bytes and QPS, not equal alphabet width. Boundaries come from sampling the actual corpus.
  • Dynamic split/merge like Bigtable tablets: a shard exceeding a size or QPS threshold splits at its median key and the registry updates.
  • Read replicas absorb read skew without resharding โ€” cheap and usually sufficient, since this workload is read-dominated.

Streaming, not materializing

A range can match arbitrarily many words, so the coordinator must stream: open an iterator per overlapping shard, k-way merge them through a min-heap, emit in sorted order, stop at the page limit. Never collect the full result set in the coordinator โ€” one range("", "") would OOM it.

cursor = base64(last_word_returned)
next page: range(cursor, b) exclusive of cursor
// stateless, resumable, survives coordinator restart

Defend it: a key-based cursor is stable under concurrent inserts, whereas an offset-based one silently skips or repeats words when the corpus changes mid-pagination.

Space: prefix compression is the one optimization worth naming

Sorted leaves mean adjacent words share long prefixes (aa, aaa, aabโ€ฆ). Store each as (shared_prefix_len, suffix) โ€” front-coding โ€” with a full key every 16 entries so you can still binary-search within a block. Typical dictionary corpora compress 3โ€“5ร— this way, which is the difference between leaves fitting in page cache and not. Cost: a decode step per entry, negligible against the SSD read it saves.

5. Follow-ups โ€” answers to have ready

How do you handle deletes during a scan?

MVCC: reads take a snapshot version at query start and see a consistent view for the whole scan, including across pages if you pin the snapshot in the cursor. Deletes write a tombstone that compaction later reclaims. Without snapshots, a long paginated scan can return a word that was deleted an hour earlier and miss one inserted before it started, which is confusing rather than merely stale.

A shard dies mid-query.

Each shard is a replica group (3 replicas, Raft or primary/backup); the coordinator retries the iterator against another replica from the last returned key, which is safe because the cursor is key-based and the scan is idempotent. The partial results already streamed to the client are still valid and in order โ€” you just resume. If a whole group is unavailable, return an explicit partial result with the missing range flagged rather than silently returning an incomplete set; silently-incomplete is the worse failure for a search-like API.

Unicode? Locale-specific collation?

"Lexicographic" is underspecified the moment you leave ASCII. Byte-order on UTF-8 is not the same as human alphabetical order in most languages, and locales genuinely disagree (in Swedish, รค sorts after z). Practical answer: store an ICU collation key alongside each word and range-partition on that, with the collation locale as part of the index identity โ€” one index per locale you support. It's a real cost and worth surfacing as a clarifying question early rather than a surprise later.

Would you use an existing system?

Yes, and saying so is the senior answer. Bigtable, HBase, and CockroachDB are all range-partitioned, sorted-key stores that give you exactly this: ordered scans, dynamic tablet splitting, replication. I'd build on one of those and spend my effort on the collation and pagination semantics. Building a distributed B+ tree from scratch is only justified if the workload has a property those don't serve โ€” and I'd want to name that property before signing up for it.

6. Numbers & recap

Sizing

  • 10ยนโฐ words ร— 20 B = 200 GB raw; ~50 GB after front-coding.
  • ~10 shards at 5 GB each, ร—3 replicas = 30 nodes. Small.
  • B+ tree fanout ~200 โ†’ depth 5 for 10ยนโฐ keys; internal nodes ~1 GB, cached in RAM, so a lookup is one SSD read.
  • Scan of 1000 results โ‰ˆ a few contiguous leaf pages โ‰ˆ < 5 ms.

30-second recap

Range queries need sorted, contiguous data, so I range-partition rather than hash-partition โ€” accepting hot shards as the explicit cost. Each shard is a replicated B+ tree with internal nodes in RAM and prefix-compressed, linked leaves on SSD, so a range is one descent plus a sequential leaf walk. A coordinator consults a boundary registry, opens iterators only on overlapping shards, and k-way merges them into a streamed, key-cursor-paginated response โ€” never materializing the full result. Skew is handled by choosing boundaries from a sample of the real corpus and splitting tablets dynamically on size or QPS, with read replicas absorbing query skew. The things I'd pin down first are the interval semantics โ€” the examples imply half-open โ€” and the collation, because lexicographic order is locale-dependent the moment the corpus isn't ASCII.

Local business search service condensed

็‹—ๅฎถ senior, 2024-09. "่ฎพ่ฎกไธ€ไธช local business search service." Given API params: geolocation (lat, lng), radius. Explicitly a senior-level prompt where the details are yours to elicit โ€” the interviewer gave the minimum and expected the candidate to drive.

1. Requirements โ€” what to elicit

Functional

  • search(lat, lng, radius, query?, filters?) โ†’ ranked businesses.
  • Filters: category, open-now, rating, price band.
  • Ranking blends distance, relevance, and quality โ€” ask which dominates; it's a product decision.
  • Business data ingest: owner edits, third-party feeds, user-suggested corrections.

Ask early: is this "find me sushi nearby" (text + geo) or "what's within 2 km" (pure geo)? The first needs an inverted index; the second only needs a spatial index. Most candidates assume one and design the wrong system.

Non-functional

  • ~10โธ businesses globally; ~10โต QPS peak, heavily skewed to dense metros.
  • p99 < 200 ms โ€” it's an interactive search box.
  • Read-dominated ~10 000:1. Business data changes slowly; open-now and busy-ness change constantly.
  • Stale results acceptable for hours on descriptions, not on permanent closure.

Core tension: geo-filtering and text-relevance want different index structures, and joining them at query time is expensive. The design is about doing the cheap filter first and the expensive scoring on a small candidate set.

2. Geo indexing โ€” the part they're actually testing

ApproachHow radius search worksTrade-off
GeohashPrefix = bounding box. Query the covering cell + 8 neighbors, then filter by exact distance.Simple, string-prefix friendly. Boundary artifacts: two nearby points can differ at the first character, hence the neighbor scan.
S2 cells (pick this)Hilbert curve โ†’ 64-bit cell IDs. A radius becomes a small set of variable-level cell ranges (S2RegionCoverer), each a contiguous ID range.Better locality, no pole/meridian distortion, and ranges map directly onto a sorted-key store. It's also what Google actually uses, which lands well here.
Quadtree / R-treeTree descent to the query rectangle.Adapts to density, but tree rebalancing under writes and awkward to shard across machines.
PostGIS / naive lat-lng boxWHERE lat BETWEEN โ€ฆ AND lng BETWEEN โ€ฆFine to ~10โถ rows on one box. At 10โธ with 10โต QPS it falls over โ€” say why, then discard it.

The move to state: S2 turns a 2-D proximity problem into 1-D range scans over sorted integers, which is something a distributed store can shard and serve. That sentence is most of the value of this section.

3. Architecture

Client lat,lng,r,q Query service S2 coverer โ†’ cell ranges Geo index shards cell_id โ†’ biz_ids sharded by cell prefix Text index inverted, per-region term โ†’ biz_ids Attribute cache hours, rating, price open-now computed Ranker L1 cheap โ†’ top 200 L2 model โ†’ top 20 Hydrate business docs from KV Ingest: owner edits / feeds / corrections โ†’ index build batch rebuild nightly + incremental delta index

Cheap filters first: S2 cell ranges cut 10โธ businesses to a few thousand candidates before anything expensive runs. Intersect with the text posting list, apply attribute filters, then two-stage rank. Hydration of full business documents happens last, on ~20 IDs.

4. Trade-offs worth stating

Density skew is the defining problem

Manhattan has orders of magnitude more businesses per kmยฒ than rural Montana, and far more queries. A fixed S2 level is therefore wrong everywhere: at level 13 a Manhattan cell holds thousands of businesses, and a Montana cell holds none.

  • Variable-level cells: subdivide until a cell holds โ‰ค K businesses. Dense areas get deep cells, sparse areas shallow ones. S2RegionCoverer handles mixed levels natively โ€” that's the reason to prefer S2 over flat geohash.
  • Adaptive radius: if a small radius in a dense area already returns hundreds of results, don't expand. If a rural query returns two, expand the radius progressively โ€” but tell the user you did, or "nearest" results 40 km away look like a bug.
  • Shard by cell prefix but balance by load, with hot metro cells split across more replicas. Geographic sharding without load balancing puts all of Tokyo on one machine.

Two-stage ranking, and why

L1 (cheap, on thousands of candidates):
  score = w1*exp(-dist/d0)
        + w2*log(1+review_count)*rating
        + w3*text_match_bm25
  โ†’ keep top 200

L2 (expensive, on 200):
  GBDT / neural with personalization,
  popularity-at-this-hour, click history
  โ†’ top 20

Running the expensive model on every candidate blows the 200 ms budget; running only the cheap one loses quality. The staged funnel is the standard answer and stating the candidate-set sizes is what makes it credible.

Defend it โ€” freshness split. Static attributes (name, category, location) go in the nightly-rebuilt index. Volatile ones (open now, temporary closure, live busy-ness) must not be baked into the index โ€” compute open-now at query time from cached hours plus the local timezone, and keep a small, fast-updating overlay for closures. Baking hours into the index means a business shows as open at 3 a.m. until the next rebuild.

5. Follow-ups โ€” answers to have ready

Why not just PostGIS with a GiST index?

For 10โถ businesses and modest QPS, that is the right answer and I'd say so โ€” don't build a distributed index you don't need. It stops working at 10โธ rows and 10โต QPS with a text-relevance join, because you can't shard a single PostGIS instance by geography without building the coordinator layer anyway, and the ranking model doesn't belong in the database. The migration path is: start with PostGIS, move the geo index to S2-over-a-sorted-store when the metro shards get hot.

A business permanently closes. How fast does it disappear?

Fast, via the delta path, not the nightly rebuild. Closure writes a tombstone into a small, always-consulted overlay that the query service intersects out of results โ€” seconds, not hours. This is the one data change where staleness is genuinely user-harmful (someone drives to a closed restaurant), so it gets its own fast path while descriptions and photos ride the batch rebuild.

Query at a cell boundary โ€” do you miss a business 10 m away in the next cell?

No, because the coverer generates the cells covering the whole circle, not the cell containing the center โ€” the covering inherently spans boundaries. Then every candidate gets an exact haversine distance check, so cell coverage only ever over-fetches, never under-fetches. That "cover generously, filter exactly" pattern is the correctness argument, and it's worth stating explicitly since boundary handling is the classic geohash bug.

Hot metro traffic melts a shard.

Read replicas first, since this is 10 000:1 read-dominated and replicas are cheap. Then a query-result cache keyed on quantized inputs โ€” snap lat/lng to a ~100 m grid and radius to buckets, so nearby users share a cache entry; that alone gives a very high hit rate in dense areas where queries cluster. Finally split the hot cells to more shards. Notably the index is nearly static, so replicating it aggressively costs storage but no consistency headache.

6. Numbers & recap

Sizing

  • 10โธ businesses ร— ~2 KB doc = 200 GB. Index postings far smaller.
  • Geo index: 10โธ (cell_id, biz_id) pairs โ‰ˆ 1.6 GB โ€” small enough to replicate widely.
  • 10โต QPS ร— ~3 KB response = 300 MB/s out. Cache-friendly.
  • Budget: coverer ~1 ms, index scan ~20 ms, L1 ~10 ms, L2 on 200 docs ~30 ms, hydrate ~10 ms โ†’ ~70 ms, well inside 200 ms.

30-second recap

The core move is turning 2-D proximity into 1-D range scans: an S2 region coverer converts (lat, lng, radius) into a handful of contiguous cell-ID ranges, which a sorted, range-partitioned index can serve directly. Cells are variable-level so a dense metro subdivides deeper than farmland, which is why S2 beats a flat geohash here. Covering is generous and then every candidate gets an exact haversine check, so boundaries can over-fetch but never miss. Text queries intersect the cell candidates with an inverted-index posting list, filters apply, then a two-stage ranker โ€” a cheap linear pass down to 200, an expensive model down to 20 โ€” keeps us inside 200 ms. Volatile attributes like open-now and permanent closures ride a fast overlay rather than the nightly index rebuild, because a user driving to a closed restaurant is the failure that actually matters. And I'd start on PostGIS if the scale were 10โถ โ€” this design only earns its complexity at 10โธ.

Quota limiter (Drive-style usage quota) condensed

From the 1p3a Google question bank, L6 System Design cluster. The bank's own framing is the point: clarify whether the requirement is request throttling, per-user/per-tenant resource quota, or storage-usage enforcement before choosing counters, windows, and reconciliation semantics. Guessing wrong here means you design a rate limiter when they wanted an accounting system.

1. Requirements โ€” the disambiguation is the first grade

Rate limitingResource quotaStorage quota (assume this)
UnitRequests per windowConcurrent slots (VMs, connections)Cumulative bytes stored
Resets?Yes, every windowOn releaseNever โ€” it's a running balance
Error of over-countSelf-heals next windowSelf-heals on releasePermanent drift โ€” needs reconciliation
Right primitiveToken bucket in RedisSemaphore / leaseDurable counter + async audit
Failure stanceFail openFail closedFail open on read, closed on write

Say all three out loud, pick one, and justify: "Drive-style implies cumulative storage, so I'll design that โ€” it's the hardest of the three because errors accumulate forever rather than washing out."

Functional (storage quota)

  • reserve(user, bytes) before an upload; commit or release after.
  • Quota is hierarchical: org โ†’ team โ†’ user, and the tightest binding limit applies.
  • Deletes and trash restore adjust usage; trash counts against quota until purged (a real product decision โ€” state it).
  • Users see accurate usage, broken down by category.

Non-functional

  • 10โน users; ~10โต quota ops/sec; p99 < 50 ms on the upload path.
  • Accuracy over availability on the write path โ€” the inverse of a rate limiter. Letting a user exceed quota costs real money and is hard to claw back.
  • Displayed usage may lag by seconds; enforced usage may not drift permanently.

Core tension: you need a durable, monotonic, per-user counter that survives crashes and never double-counts, on a path that also needs to be fast. That's a transaction, not a cache.

2. Design โ€” reserve / commit with reconciliation

The two-phase sequence

reserve(user, bytes, op_id):
  BEGIN
    row = SELECT used, reserved, limit
            FROM quota WHERE user=? FOR UPDATE
    if row.used + row.reserved + bytes > row.limit:
        ROLLBACK; return 507 Insufficient Storage
    INSERT INTO reservation(op_id, user, bytes,
              expires_at=now+1h)
         ON CONFLICT (op_id) DO NOTHING   -- idempotent
    UPDATE quota SET reserved = reserved + bytes
  COMMIT

commit(op_id):   -- after bytes are durable in blob store
  BEGIN
    r = DELETE FROM reservation WHERE op_id=? RETURNING *
    if r is null: return OK          -- already committed
    UPDATE quota SET used = used + r.bytes,
                     reserved = reserved - r.bytes
  COMMIT

sweeper (every 5 min):
  expire reservations past expires_at -> release

Why reserve-then-commit, not just increment

A single increment at upload start over-counts every abandoned upload; a single increment at upload end lets a user start a thousand concurrent uploads that each individually fit under the limit and blow past it together. The reservation makes the check and the claim atomic, which is the actual race in this problem.

op_id makes both phases idempotent, so a client retry after a timeout can't double-charge. The expiry sweeper is what stops a crashed uploader from permanently stranding quota โ€” without it, reserved-but-never-committed bytes leak until the user is locked out of their own account with no way to recover.

Why a transactional store (Spanner/SQL) rather than Redis: this counter is the source of truth for something users pay for. It must survive a cache flush, support a real transaction across the reservation and the balance, and be auditable. This is the exact inverse of the rate-limiter argument in panel 1 โ€” and being able to explain why the same-looking problem gets the opposite answer is the whole point of having both.

3. Trade-offs worth stating

Hot row on shared quotas

An org-level quota is one row that every member's upload contends on with FOR UPDATE. A 10 000-person org serializes on it.

  • Sharded counters: split the org balance into N sub-rows; a reserve picks one at random, and only when a sub-row is exhausted does it consult siblings. Reads sum all N. Cost: the "am I over?" check is approximate near the boundary, so keep a small reserve buffer.
  • Lease blocks: a service instance leases 1 GB of org quota and hands out sub-allocations locally. Same pattern as the rate limiter's token lease, but here the lease must be durable and reclaimable, not best-effort.

Drift, and the audit job

Because errors are permanent, you need a periodic reconciler that recomputes true usage from the authoritative object metadata (sum of object sizes per owner) and compares it to the counter. Discrepancies get logged, and small ones auto-corrected; large ones page, because a systematic drift means a bug that's silently mischarging users.

Defend it: the reconciler is scanning billions of objects, so run it incrementally โ€” per user, on a rolling schedule, prioritized by users near their limit and by accounts with recent anomalies. A full-fleet nightly recompute doesn't scale and isn't necessary; what matters is that every account gets reconciled on some bounded horizon.

4. Follow-ups โ€” answers to have ready

Quota service is down. Can users upload?

Split by direction. Reads and display of usage fail open with a cached value โ€” showing a slightly stale number is harmless. Writes fail closed for users near their limit and open for users far below it: if the last known usage is under, say, 80%, admit the upload and reconcile later, because the worst case is a small, correctable overage. If they're near the limit, reject with a retryable error. That graded stance is much better than a blanket answer, and it's exactly the kind of "which is right for this system" reasoning the L6 rubric asks for.

A user deletes 10 GB. When does their quota free up?

Depends on the product decision about trash, which you should surface rather than assume: if deleted files sit in trash for 30 days and remain restorable, they must still count, otherwise you've promised storage you can't guarantee on restore. So the counter decrements at purge, not at delete, and the UI must say so โ€” this is the single most common user complaint about storage quotas and naming it shows product sense.

Deduplication โ€” two users store the same file. Who pays?

Both, for their logical usage, even though you store one copy. Charging by physical bytes would leak information (usage dropping tells you someone else has the same file) and produce a quota that changes without the user doing anything. So logical accounting for quota, physical accounting for cost โ€” and keeping those two ledgers separate is the right architecture.

Why is this different from the rate limiter you designed earlier?

Because the error term behaves differently. A rate limiter's over-count vanishes at the next window, so you can trade accuracy for latency with local leases and Redis, and fail open. A storage quota's over-count is permanent and monetary, so it needs a durable transactional counter, idempotent two-phase updates, an expiry sweeper, and a reconciliation job โ€” and it fails closed at the boundary. Same shape, opposite answers, and the reason is entirely about whether the error self-heals.

5. Numbers & recap

Sizing

  • 10โน users ร— ~200 B quota row = 200 GB. Sharded SQL/Spanner, easily.
  • 10โต quota ops/s, each a short transaction on one row โ†’ ~2 000 rows/s/shard across 50 shards. Comfortable.
  • Reservations in flight: 10โต ops/s ร— ~60 s mean upload = 6 M open rows. Small, and the sweeper keeps it bounded.
  • Reconciler: 10โน users on a 7-day rolling horizon โ‰ˆ 1 650 accounts/s of background scan.

30-second recap

First I'd pin down which quota this is, because request throttling, concurrency slots, and cumulative storage need genuinely different machinery. Assuming Drive-style storage: the counter is a durable, transactional per-user balance, not a cache, because an over-count here is permanent and monetary rather than washing out at the next window. Uploads take a reserve-then-commit path โ€” one transaction atomically checks the limit and claims the bytes, keyed by an operation ID so retries are idempotent, and an expiry sweeper releases reservations from crashed uploads so quota can't leak. Shared org quotas would serialize on one hot row, so I'd shard the balance into sub-counters or hand out durable lease blocks, accepting approximate enforcement near the boundary in exchange for concurrency. Because errors accumulate, there's a rolling reconciler that recomputes true usage from object metadata and pages on systematic drift. And the failure stance is graded rather than binary: fail open for users well below their limit, closed for users near it.

Distributed key-value store + throughput estimation condensed

Reported in an earlier Google SD round (the "System Design ๆŒ‚็ป" thread): design a distributed KV system, with the interviewer pushing on throughput estimation. That second half is the tell โ€” this round is at least as much about capacity math out loud as about architecture.

1. Requirements

Functional

  • get(key), put(key, value), delete(key); optionally cas(key, value, expected_version).
  • Keys โ‰ค 256 B, values โ‰ค 1 MB.
  • Optional TTL per key.
  • Ask: do we need range scans? If yes the partitioning answer flips from hash to range โ€” see panel 7.

Non-functional โ€” pin these numerically

  • 10โถ QPS, 90/10 read/write.
  • p99 < 10 ms read, < 20 ms write.
  • 10 TB of data, growing 2ร—/year.
  • Durability: no acknowledged write may be lost. Availability target 99.99%.
  • Consistency: ask, don't assume. Linearizable, read-your-writes, or eventual? This single answer determines the entire replication design.

Core tension: the CAP choice isn't philosophical here, it's a product question โ€” a session store can be eventually consistent and stay up through a partition; a counter or a lock service cannot. Make the interviewer pick, then design decisively.

2. The throughput math โ€” do this out loud, early

From QPS to node count

Reads:  9e5 QPS.  Writes: 1e5 QPS.

Per-node capability (commodity, NVMe):
  memory hit      ~200k ops/s
  NVMe random read ~100k IOPS, ~100 us
  network          10 Gbps = 1.25 GB/s

Data: 10 TB.  Cache budget: 10% hot = 1 TB RAM.
  at 128 GB/node -> 8 nodes just to hold cache
  say 20 nodes -> 500 GB data + 50 GB RAM each

Read path:
  90% cache hit  -> 8.1e5 ops/s from RAM
  10% disk       -> 9e4 IOPS spread over 20 nodes
                 = 4.5k IOPS/node  (of ~100k) OK

Write path with RF=3, quorum W=2:
  1e5 client writes -> 3e5 physical writes
  = 15k writes/s/node, ~1 KB each = 15 MB/s
  plus WAL fsync: batch-commit groups of ~100
  -> 150 fsync/s/node, trivially fine

Network per node:
  (8.1e5/20)*1KB read + replication traffic
  ~ 40 MB/s + 45 MB/s = 85 MB/s of 1250 MB/s
  -> network is NOT the bottleneck; RAM is.

What the math is for

The conclusion โ€” RAM for the working set, not IOPS or network, sets the node count โ€” is the thing to say. Any candidate can propose consistent hashing; few can tell you which resource actually binds, and that's what "throughput estimation" is probing.

State the assumptions as you go and keep them round: 100k IOPS, 200k memory ops/s, 10 Gbps, 1 KB average value. Nobody is checking your arithmetic to two decimals; they're checking that you know which numbers matter and roughly how big they are.

Then use it. The math should change a decision: here it says 20 nodes is right, that we're memory-bound so a bigger cache tier is the cheapest lever, and that RF=3 quorum writes cost us nothing we can't afford. Math that doesn't change a decision is decoration.

3. Architecture

Client lib caches ring Coordinator any node can be hash(key) โ†’ vnodes Replica A (leader) WAL โ†’ memtable โ†’ SSTables Replica B other rack Replica C other zone Storage engine LSM + bloom filter block cache Membership gossip + version Repair hinted handoff, merkle Quorum N=3, W=2, R=2 W+R > N โ†’ read sees latest ack'd write tune per call site

Consistent hashing with virtual nodes (~256 per physical node) so adding a machine moves ~1/N of the data instead of half the ring, and so a heterogeneous fleet can be weighted by capacity. Replication factor 3 placed across racks and zones โ€” replica placement is where availability actually comes from, not the replication factor itself.

4. Trade-offs worth stating

Quorum, and the honest caveat

W + R > N guarantees the read set intersects the write set, so a read sees the latest acknowledged write. With N=3: W=2/R=2 is the balanced default; W=3/R=1 favors read-heavy workloads at the cost of write availability; W=1/R=1 is fast and eventually consistent.

Defend it: quorum is not linearizability. Concurrent writers can still produce conflicting versions, a failed write that reached one replica may later surface, and read-repair timing is not deterministic. If the interviewer wants true linearizability, say so plainly and switch to Raft/Paxos per partition with reads served by the leader (or via lease reads) โ€” that costs you a leader election window on failover, which quorum-only designs don't have. Naming this distinction is a strong senior signal; asserting "quorum gives strong consistency" is a common and visible error.

Conflicts, hot keys, LSM cost

  • Conflict resolution: last-write-wins on a timestamp is simple and silently loses data under clock skew. Vector clocks / version vectors preserve causality but push merge logic to the client. State the choice โ€” LWW is defensible for a cache-like store, unacceptable for a shopping cart.
  • Hot key: hashing spreads keys, not traffic to one key. One viral key still lands on 3 replicas. Fixes: client-side caching with short TTL, or replicating that key to extra nodes on demand. Say that consistent hashing doesn't solve this โ€” many candidates think it does.
  • LSM trade: writes are sequential and fast; reads may touch multiple levels (mitigated by bloom filters), and compaction consumes background IO and causes p99 spikes. If the workload were read-dominated with in-place updates, a B-tree would be the better pick โ€” and knowing when not to use an LSM is the point.

5. Follow-ups โ€” answers to have ready

A node dies. What happens to writes targeting it?

With W=2 of N=3, writes still succeed on the surviving two โ€” availability is preserved by construction. The coordinator stores a hinted handoff for the dead replica and replays it when the node returns. If the node is gone for longer than the hint window, an anti-entropy repair using Merkle trees reconciles the ranges by exchanging hashes rather than data, so only the differing subranges are shipped. The cost of getting this wrong is silent permanent divergence, which is why the repair job is not optional.

How do you add a node without a latency spike?

Virtual nodes mean the new machine claims ~1/N of the ranges from many existing nodes rather than a contiguous half from one neighbor, so the migration is spread and no single donor saturates. Stream ranges in the background at a throttled rate, serve reads from the old owner until a range is fully transferred, then flip ownership atomically in the membership version. The throttle is the important knob โ€” an unthrottled rebalance is itself an outage, and I'd bound it to a fraction of each node's disk and network budget.

Read-your-writes for a specific user?

Three options, cheapest first: route that user's requests to the same coordinator and replica set via sticky routing; or have the client carry the version it last wrote and require the read to see at least that version, retrying elsewhere if not; or use W=3. The middle option โ€” a monotonic version token in the client โ€” is the general and correct one, and it's the same mechanism as a session consistency token in Spanner or DynamoDB.

Where does this design stop working?

At 2ร— growth per year, the memory-bound conclusion holds until the hot working set stops fitting economically in RAM โ€” around 10ร— current size. At that point the lever isn't more nodes, it's a tiered cache with a dedicated hot tier, or accepting a lower cache hit rate and re-doing the IOPS math. The other break point is multi-region: everything above assumes one region, and crossing regions forces an explicit choice between async replication with conflict handling and synchronous quorums at ~100 ms. I'd make that a separate design conversation rather than hand-waving it.

6. Numbers & recap

The headline figures

  • 10 TB ร— RF 3 = 30 TB physical; 20 nodes ร— ~2 TB NVMe. Comfortable.
  • Hot set 10% = 1 TB RAM across the fleet โ†’ 50 GB/node โ€” this is the binding constraint.
  • Reads: 9e5 QPS, ~90% from memory; disk load only ~4.5k IOPS/node of ~100k available.
  • Writes: 1e5 ร— RF 3 = 3e5/s = 15k/s/node โ‰ˆ 15 MB/s, with batched fsync at ~150/s.
  • Network ~85 MB/s/node against 1.25 GB/s โ€” an order of magnitude of headroom.

30-second recap

Consistent hashing with a few hundred virtual nodes per machine spreads both data and rebalancing cost, and replication factor 3 placed across racks and zones is where availability actually comes from. Each node runs an LSM engine โ€” write-ahead log, memtable, SSTables with bloom filters โ€” which suits a write-heavy path at the cost of read amplification and compaction-driven p99 spikes. Quorum with W=2, R=2 of N=3 means a read intersects the latest acknowledged write, but I'd be explicit that this is not linearizability: for that you'd want Raft per partition with leader reads, trading a failover window for stronger guarantees. Failed nodes are covered by hinted handoff and Merkle-tree anti-entropy repair. On the sizing: at a million QPS with 10 TB and a 10% hot set, it's memory that sets the node count at around twenty โ€” disk IOPS and network both have an order of magnitude of headroom โ€” so the cheapest lever if we need more throughput is a bigger cache tier, not more machines.